【问题标题】:Akka Streams takeWhile processing next element even after condition fails即使条件失败,Akka Streams takeWhile 也会处理下一个元素
【发布时间】:2020-11-19 05:09:23
【问题描述】:

我有一个非常简单的演员,它只打印数字:-

public class PrintLineActor extends AbstractLoggingActor {

  @Override
  public Receive createReceive() {
    return receiveBuilder()
        .match(Integer.class, i -> {
          System.out.println("Processing: " + i);
          sender().tell(i, self());
        }).build();
  }
}

现在,我有一个流来打印偶数,直到遇到奇数元素:-

  @Test
  public void streamsTest() throws Exception {

    ActorSystem system = ActorSystem.create("testSystem");
    ActorRef printActor = system.actorOf(Props.create(PrintLineActor.class));

    Integer[] intArray = new Integer[]{2,4,6,8,9,10,12};
    CompletionStage<List<Integer>> result = Source.from(Arrays.asList(intArray))
        .ask(1, printActor, Integer.class, Timeout.apply(10, TimeUnit.SECONDS))
        .takeWhile(i -> i != 9)
        .runWith(Sink.seq(), ActorMaterializer.create(system));

    List<Integer> result1 = result.toCompletableFuture().get();
    System.out.println("Result :- ");
    result1.forEach(System.out::println);
  }

不期望处理 9 之后的任何元素,也就是发送给 actor。但是,我看到演员(但不是 12)也在处理数字“10”,如下面的输出所示

Processing: 2
Processing: 4
Processing: 6
Processing: 8
Processing: 9
Processing: 10 //WHY IS THIS BEING PROCESSED BY ACTOR??

Result :- 
2
4
6
8

为什么演员要处理 10?如何阻止这种情况?

编辑:

我已经尝试通过记录事件的时间戳来进行调试,只是为了查看是否在 9 实际完成之前处理了 10,但是没有,在完全处理 9 之后取 10。这是日志:-

Before Ask: 2 in 1596035906509
Processing inside Actor: 2 at 1596035906509
Inside TakeWhile 2 at  in 1596035906509

Before Ask: 4 in 1596035906609
Processing inside Actor: 4 at 1596035906610
Inside TakeWhile 4 at  in 1596035906610

Before Ask: 6 in 1596035906712
Processing inside Actor: 6 at 1596035906712
Inside TakeWhile 6 at  in 1596035906712

Before Ask: 8 in 1596035906814
Processing inside Actor: 8 at 1596035906814
Inside TakeWhile 8 at  in 1596035906815

Before Ask: 9 in 1596035906915
Processing inside Actor: 9 at 1596035906915
Inside TakeWhile 9 at  in 1596035906916

Before Ask: 10 in 1596035907017 //so 10 is taken much after the 9 is processed fully
Processing inside Actor: 10 at 1596035907017

Result :- 
2
4
6
8

此外,如果我将 .ask 替换为直接 .map(print..),则不会打印 10。所以当涉及到 actor.ask 时为什么会发生这种情况对我来说很奇怪。

【问题讨论】:

    标签: java scala akka akka-stream


    【解决方案1】:

    因为你ask printActor 是异步而不是同步打印。在你的确切运行中:

    • 消息 9 到达 PrintActor,打印“Processing: 9”
    • 消息 10 到达 PrintActor,打印“Processing: 10”
    • 您的 Akka 流从 PrintActor 接收到消息 9 的响应消息,完成 Akka 流,因此结果中既没有 9 也没有 10。

    要解决确切的问题,请删除异步询问并改为同步打印。但不确定 PrintActor 是否只是一个类比,请告诉我。

    【讨论】:

    • 不,这不是事件的顺序。此流是顺序的,没有并行性。我用一些调试日志编辑了这个问题,以指示事件的顺序。 9处理成功后要求actor处理10。
    • 我不是说,9 晚于 10 处理。我是说在 PrintActor 处理 10 之后,您的 Akka 流会收到消息 9 的询问响应,并在收到来自消息 9 的询问响应之前完成流程消息 10。这就是原因。
    • 不,消息 9 的询问响应在消息 10 发送给 actor 之前出现。您可以在调试日志中看到这一点
    • 请更新有关打印日志位置的代码
    • 嘿,我在 .ask 之前放了一个 map(x -> {print("Before ask: "+ x); return x;) 我还在 takeWhile 中打印了时间戳。看起来 mapAsync/ask 在发出元素后立即从上游拉取元素。所以,就像你说的,没有隐含的方法可以做到这一点:(
    【解决方案2】:

    akka 流通过不同的流缓冲值。请注意,10 已被处理,但它不是结果的一部分。如果您愿意,可以配置缓冲区大小:

    .ask(1, printActor, Integer.class, Timeout.apply(10, TimeUnit.SECONDS)).buffer(1, OverflowStrategy.backpressure)
    

    【讨论】:

    • 感谢您的回复。这不起作用 - 10 仍然由演员处理。我尝试在询问之前放入 .buffer(1,..),在询问之后,演员仍然以某种方式得到 10
    • 另外,我不希望处理 10 个
    • 不幸的是,您无法控制 akka 流停止的方式。这是你能做的最好的。为了更好地理解这一点,您可以在接收器中打印带有时间戳的结果,并查看之前发生了哪一个。重复几次。我敢肯定,在 9 被拒绝之前,至少会处理一次 10。
    • 我还在takeWhile()中添加了100ms的延迟。所以在 takewhile() 之后的 100ms 延迟后,10 从源中取出
    猜你喜欢
    • 2014-12-03
    • 1970-01-01
    • 1970-01-01
    • 2016-06-23
    • 2019-12-13
    • 1970-01-01
    • 2011-02-05
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多