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
相关产品推荐
相关产品推荐

