【问题标题】:Akka Streams Accumulate by Source SingleAkka Streams 按源单累积
【发布时间】:2020-09-07 15:22:06
【问题描述】:

我正在尝试使用 akka 流来积累数据并用作批处理:

val myFlow: Flow[String, Unit, NotUsed] = Flow[String].collect {
    case record =>
      println(record)
      Future(record)
  }.mapAsync(1)(x => x).groupedWithin(3, 30 seconds)
    .mapAsync(10)(records =>
      someBatchOperation(records))
    )

我对上面代码的期望是直到 3 条记录准备好或 30 秒过去后才进行任何操作。但是当我用Source.single("test") 发送一些请求时,它正在处理这条记录,而无需等待其他人或30秒。

如何使用此流程来等待其他记录到来或 30 秒空闲?

记录来自一个 API 请求,我正在尝试在流中累积这些数据,例如:

Source.single(apiRecord).via(myFlow).runWith(Sink.ignore)

【问题讨论】:

    标签: scala akka akka-stream


    【解决方案1】:

    它确实做到了。让我们考虑以下几点:

    Source(Stream.from(1)).throttle(1, 400 milli).groupedWithin(3, 1 seconds).runWith(Sink.foreach(i => println(s"Done with ${i} ${System.currentTimeMillis}")))
    

    在我终止进程之前,该行的输出是:

    Done with Vector(1, 2, 3) 1599495716345
    Done with Vector(4, 5) 1599495717348
    Done with Vector(6, 7, 8) 1599495718330
    Done with Vector(9, 10) 1599495719350
    Done with Vector(11, 12, 13) 1599495720330
    Done with Vector(14, 15) 1599495721350
    Done with Vector(16, 17, 18) 1599495722328
    Done with Vector(19, 20) 1599495723348
    Done with Vector(21, 22, 23) 1599495724330
    

    正如我们所见,每次发射 2 个元素到 3 个元素之间的时间差略大于 1 秒。这是有道理的,因为在 1 秒延迟之后,到达打印线需要更多时间。

    我们每次发射 2 个元素到 3 个元素的时间差不到一秒。因为它有足够的元素继续下去。

    为什么它在您的示例中不起作用?

    当您使用Source.single 时,源会为其自身添加一个完整的阶段。您可以在source code of akka 中看到它。 在这种情况下,groupedWithin 流知道它不会再获得任何元素,因此它可以发出"test" 字符串。为了实际测试这个流,尝试创建一个更大的流。

    使用 Source(1 到 10) 时,它实际上转换为 Source.Single,这也完成了流。我们可以看到here

    【讨论】:

    • 但是去掉油门后,流程会立即完成。你能解释一下为什么吗?
    • @YikSanChan 如果我们取消油门,那么我们只会达到 3 个数字的限制,这将一直发生。油门只是在这里帮助我们看到两个阈值都被击中。 400 毫秒是一个时间量,有时 3 可以在同一秒内发生,有时不能
    • 感谢您的努力,但节流并不能直接挽救我的问题。我的记录一一来自api,我正在将传入的数据发送到流中。我想累积这些记录,直到它们的计数超过或时间到了,但是当我使用 source.single 时,它并没有在流程中的某个地方累积。
    • 正如我的帖子中所解释的,这是不可能的。我再重复一遍:当使用Source.single 时,将完成阶段添加到源中。流程知道它不会再获得任何元素,所以它会发出它拥有的东西。此外,发送 1 个元素并试图累积 3 个元素是没有意义的。你怎么能期待呢?你应该有一个来源,接收你所有的元素,而不是积累这些。
    猜你喜欢
    • 2017-11-27
    • 2019-04-13
    • 1970-01-01
    • 1970-01-01
    • 2018-04-02
    • 1970-01-01
    • 2017-06-07
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多