【问题标题】:Extract values from list of tuples continuously in Erlang在 Erlang 中连续从元组列表中提取值
【发布时间】:2018-05-25 17:28:58
【问题描述】:

我正在从 Ruby 的背景中学习 Erlang,并且在掌握思维过程方面有些困难。我要解决的问题如下:

我需要向 api 发出相同的请求,每次我在响应中收到一个唯一 ID,我需要将其传递给下一个请求,直到没有返回 ID。从每个响应中,我都需要提取某些数据并将其用于其他事情。

首先获取迭代器:

ShardIteratorResponse = kinetic:get_shard_iterator(GetShardIteratorPayload).
{ok,[{<<"ShardIterator">>,
      <<"AAAAAAAAAAGU+v0fDvpmu/02z5Q5OJZhPo/tU7fjftFF/H9M7J9niRJB8MIZiB9E1ntZGL90dIj3TW6MUWMUX67NEj4GO89D"...>>}]}

解析出 shard_iterator..

{_, [{_, ShardIterator}]} = ShardIteratorResponse.

向 kinesis 请求流记录...

GetRecordsPayload = [{<<"ShardIterator">>, <<ShardIterator/binary>>}].
[{<<"ShardIterator">>,
  <<"AAAAAAAAAAGU+v0fDvpmu/02z5Q5OJZhPo/tU7fjftFF/H9M7J9niRJB8MIZiB9E1ntZGL90dIj3TW6MUWMUX67NEj4GO89DETABlwVV"...>>}]
14> RecordsResponse = kinetic:get_records(GetRecordsPayload).
{ok,[{<<"NextShardIterator">>,
      <<"AAAAAAAAAAFy3dnTJYkWr3gq0CGo3hkj1t47ccUS10f5nADQXWkBZaJvVgTMcY+nZ9p4AZCdUYVmr3dmygWjcMdugHLQEg6x"...>>},
     {<<"Records">>,
      [{[{<<"Data">>,<<"Zmlyc3QgcmVjb3JkISEh">>},
         {<<"PartitionKey">>,<<"BlanePartitionKey">>},
         {<<"SequenceNumber">>,
          <<"49545722516689138064543799042897648239478878787235479554">>}]}]}]}

我正在苦苦挣扎的是如何编写一个循环,该循环不断地到达该流的 kinesis 端点,直到没有更多的分片迭代器,也就是我想要所有记录。因为我不能像在 Ruby 中那样重新分配变量。

【问题讨论】:

    标签: erlang


    【解决方案1】:

    警告:我的代码可能有问题,但它是“关闭的”。我从未运行过它,也看不到最后一个迭代器的样子。

    我看到您正在尝试完全在 shell 中完成您的工作。这是可能的,但很难。可以使用命名函数和递归(since release 17.0 it's easier),例如:

    F = fun (ShardIteratorPayload) ->
        {_, [{_, ShardIterator}]} = kinetic:get_shard_iterator(ShardIteratorPayload),
        FunLoop =
            fun Loop(<<>>, Accumulator) ->  % no clue how last iterator can look like
                    lists:reverse(Accumulator);
                Loop(ShardIterator, Accumulator) ->
                    {ok, [{_, NextShardIterator}, {<<"Records">>, Records}]} =
                        kinetic:get_records([{<<"ShardIterator">>, <<ShardIterator/binary>>}]),
                    Loop(NextShardIterator, [Records | Accumulator])
            end,
        FunLoop(ShardIterator, [])
    end.
    AllRecords = F(GetShardIteratorPayload).
    

    但是在shell中打字太复杂了……

    在模块中编写代码要容易得多。 erlang 中的一个常见模式是生成另一个或多个进程来获取数据。为简单起见,您可以通过调用 spawn or spawn_link 来生成另一个进程,但现在不要打扰链接,只需使用 spawn/3。 让我们编译简单的消费者模块:

    -module(kinetic_simple_consumer).
    
    -export([start/1]).
    
    start(GetShardIteratorPayload) ->
        Pid = spawn(kinetic_simple_fetcher, start, [self(), GetShardIteratorPayload]),
        consumer_loop(Pid).
    
    consumer_loop(FetcherPid) ->
        receive
            {FetcherPid, finished} ->
                ok;
            {FetcherPid, {records, Records}} ->
                consume(Records),
                consumer_loop(FetcherPid);
            UnexpectedMsg -> 
                io:format("DROPPING:~n~p~n", [UnexpectedMsg]),
                consumer_loop(FetcherPid)
        end.
    
    consume(Records) ->
        io:format("RECEIVED:~n~p~n",[Records]).
    

    还有抓取器:

    -module(kinetic_simple_fetcher).
    
    -export([start/2]).
    
    start(ConsumerPid, GetShardIteratorPayload) ->
        {ok, [ShardIterator]} = kinetic:get_shard_iterator(GetShardIteratorPayload),
        fetcher_loop(ConsumerPid, ShardIterator).
    
    fetcher_loop(ConsumerPid, {_, <<>>}) -> % no clue how last iterator can look like
        ConsumerPid ! {self(), finished};
    
    fetcher_loop(ConsumerPid, ShardIterator) ->
        {ok, [NextShardIterator, {<<"Records">>, Records}]} = 
            kinetic:get_records(shard_iterator(ShardIterator)),
        ConsumerPid ! {self(), {records, Records}},
        fetcher_loop(ConsumerPid, NextShardIterator).
    
    shard_iterator({_, ShardIterator}) ->
        [{<<"ShardIterator">>, <<ShardIterator/binary>>}].
    

    正如您所见,两个进程可以同时完成它们的工作。 从你的 shell 中尝试:

    kinetic_simple_consumer:start(GetShardIteratorPayload).
    

    现在您看到您的 shell 进程转向消费者,并且在 fetcher 发送 {ItsPid, finished} 后您将恢复您的 shell。

    下次代替

    kinetic_simple_consumer:start(GetShardIteratorPayload).
    

    运行:

    spawn(kinetic_simple_consumer, start, [GetShardIteratorPayload]).
    

    您应该使用生成过程 - 这是 erlang 的主要优势。

    【讨论】:

      【解决方案2】:

      在 Erlang 中,您可以使用尾递归函数编写循环。我不知道动力学 API,所以为了简单起见,我只是假设,当没有更多分片时,kinetic:next_iterator/1 返回{ok, NextIterator}{error, Reason}

      loop({error, Reason}) ->
          ok;
      loop({ok, Iterator}) ->
          do_something_with(Iterator),
          Result = kinetic:next_iterator(Iterator),
          loop(Result).
      

      您正在用迭代替换循环。第一个子句处理没有更多分片的情况(总是以结束条件开始递归)。第二个子句处理 case,我们得到了一些迭代器,我们用它做一些事情并调用 next。

      递归调用是函数体中的最后一条指令,称为尾递归。 Erlang 优化了此类调用 - 它们不使用调用堆栈,因此它们可以在常量内存中无限运行(您不会得到类似“堆栈级别太深”之类的东西)

      【讨论】:

        猜你喜欢
        • 2018-04-07
        • 2020-10-04
        • 1970-01-01
        • 2014-06-30
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2016-07-12
        相关资源
        最近更新 更多