You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 09:27:04