【发布时间】: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