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

Akka Stream拼接两个Flow时输出不符合预期的问题排查与实现修正

问题根因

你当前的流拓扑逻辑完全不符合需求,且触发了Akka Stream的背压机制导致流停滞,具体原因如下:

  • 需求是「所有源元素先全部经过flow1处理输出,再全部经过flow2处理输出」,但你使用的Broadcast组件会把每一个输入元素同时分发给所有下游,和需求逻辑本身就不匹配。
  • Broadcast默认启用背压对齐规则:只要任意一个下游端口处于背压状态,它就会停止发送所有元素。而Concat组件只有当第一个输入端口的流完全终止后,才会开始消费第二个输入端口的数据,第二个端口在第一个端口处理完成前会一直返回背压信号。
  • 第一个元素输入后,flow1处理后的结果通过Concat第一个端口输出打印,此时Broadcast尝试给flow2的下游发送数据,但Concat第二个端口正处于背压状态,导致Broadcast被阻塞,整个流不再处理后续元素,因此只会打印第一个结果2。

修复方案

方案1:直接拼接两个处理后的源(推荐)

不需要自定义GraphDSL,直接把源分别应用两个Flow后用concat拼接即可:

val source = Source(1 to 10)
val sink = Sink.foreach(println)

val flow1 = Flow[Int].map(_ + 1)
val flow2 = Flow[Int].map(_ * 10)

val resultSource = Source.concat(
  source.via(flow1),
  source.via(flow2)
)

resultSource.runWith(sink)

方案2:保持单输入Flow的结构

如果业务需要必须封装成单输入的自定义Flow,可以用prefixAndTail拿到完整流后物化两次再拼接:

val source = Source(1 to 10)
val sink = Sink.foreach(println)

val flow1 = Flow[Int].map(_ + 1)
val flow2 = Flow[Int].map(_ * 10)

val flowGraph = Flow[Int].prefixAndTail(0).flatMapConcat { case (_, rest) =>
  Source.concat(
    rest.via(flow1),
    rest.via(flow2)
  )
}

source.via(flowGraph).runWith(sink)

内容的提问来源于stack exchange,提问作者devj

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 07:27:03