如何基于Source或Flow中的元素计算生成对应物化值
实现方案
你需要构造一个透传所有元素、同时累计元素值作为物化值的Flow,核心利用Akka Streams的alsoToMat组合累加Sink实现,完全不需要修改你给出的固定运行代码。
完整代码实现
生产环境推荐用法(物化值为Future[Int],无阻塞)
import akka.stream.scaladsl.{Flow, Sink} import scala.concurrent.Future val dumbFlow: Flow[Int, Int, Future[Int]] = // 透传所有元素的同时,把元素副本发送给累加Sink Flow[Int].alsoToMat(Sink.fold[Int](0)(_ + _))(Keep.right)
此时运行你给出的固定代码后,result类型为Future[Int],可通过回调、map等异步方式获取求和结果6,不会阻塞线程。
完全对齐示例的实现(物化值为Int,仅测试用)
如果需要和你示例中直接返回Int的行为对齐,可额外加阻塞读取(生产环境禁止使用,会阻塞线程):
import akka.stream.scaladsl.{Flow, Sink} import scala.concurrent.Await import scala.concurrent.duration._ val dumbFlow: Flow[Int, Int, Int] = Flow[Int].alsoToMat(Sink.fold[Int](0)(_ + _))(Keep.right) // 阻塞等待流运行完成,获取累加结果,超时时间可根据流长度调整 .mapMaterializedValue(fut => Await.result(fut, 1.second))
运行你给出的固定代码后,result的值直接为6,和预期一致。
实现原理
alsoToMat是Akka Streams提供的分支算子:主流程的所有元素会原样透传给下游,同时会把每个元素的副本发送给传入的侧链Sink,不会影响主流程的运行逻辑,因此下游的Sink.ignore可以正常消费流直到结束。Sink.fold[Int](0)(_ + _)作为侧链Sink负责累加所有收到的元素,它的物化值是Future[Int]类型,当流正常完成时,Future会填充为所有元素的累加和。Keep.right表示我们选择侧链累加Sink的物化值,作为整个dumbFlow的物化值,满足你把计算逻辑放在Source侧的要求。- 如果需要基于元素做其他计算(比如计数、求最大值、自定义聚合),只需要替换
Sink.fold的初始值和聚合逻辑即可,方案通用。
内容的提问来源于stack exchange,提问作者Random42
相关产品推荐
相关产品推荐

