Akka Streams:基于Flow构建Graph时如何创建Inlet和Outlet
问题修复方案
修正后的完整代码如下:
import akka.stream.scaladsl.{Balance, Flow, GraphDSL, ZipWith} import akka.stream.OverflowStrategy import akka.stream.scaladsl.FlowShape import akka.NotUsed object Test extends App { val SomeComplicatedFlow: Flow[Int, Int, NotUsed] = Flow.fromGraph(GraphDSL.create() { implicit builder => import GraphDSL.Implicits._ // 显式将buffer注册到构建器,直接拿它的输入端口作为整个Flow的入口 val buffer = builder.add(Flow[Int].buffer(12, OverflowStrategy.backpressure)) val balance = builder.add(Balance[Int](2)) val flow1 = Flow[Int].map(_ * 2) val flow2 = Flow[Int].map(_ * 2) val zip = builder.add(ZipWith[Int, Int, Int]((left, right) => left + right)) // 显式将flow3注册到构建器,拿它的输出端口作为整个Flow的出口 val flow3 = builder.add(Flow[Int].map(_ * 2)) buffer ~> balance.in balance.out(0) ~> flow1 ~> zip.in0 balance.out(1) ~> flow2 ~> zip.in1 zip.out ~> flow3 FlowShape(buffer.in, flow3.out) }) }
两个问题的具体解决逻辑
- 问题1:你不需要额外新增单输出的Balance组件做入口,只需要把首个
buffer组件通过builder.add()注册到图构建器,就能直接获取它的Inlet作为整个自定义Flow的输入端口,完全消除冗余组件。 - 问题2:你无法直接从
flow3获取Outlet的原因是你没有将它显式注册到构建器,只要调用builder.add(flow3)拿到对应的FlowShape实例,就能直接取它的out端口作为整个Flow的出口,不需要额外新增Balance做中转。
补充:你代码里的
flow1、flow2没有显式调用builder.add也能正常连接,是因为~>运算符会隐式将未注册的流组件注册到当前构建器,只是你不需要获取它们的端口所以不需要显式处理。
内容的提问来源于stack exchange,提问作者TomatoFarmer
相关产品推荐
相关产品推荐

