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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 18:24:06