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

Akka将Flow转换为集合或Publisher 拆分Source为独立流的实现问题

实现方案

你原有代码中已经完成了流的分区逻辑,只需要调整分区方法的物化值传递逻辑,就能直接拿到你需要的两类结果:

第一步:调整扩展方法,保留Sink的物化值

原来的partition方法固定返回NotUsed,会丢弃两个传入Sink的物化结果,修改后支持带出两个Sink的返回值:

implicit class EitherSourceExtension[L, R, Mat](source: Source[FormData.BodyPart, Mat]) {
  // 新增泛型参数承接两个Sink的物化值类型
  def partition[MatL, MatR](left: Sink[BodyPartEntity, MatL], right: Sink[BodyPartEntity, MatR]): Graph[ClosedShape, (MatL, MatR)] = {
    // 传入两个Sink到GraphDSL,组合二者的物化值作为整体返回
    GraphDSL.create(left, right)((leftMat, rightMat) => (leftMat, rightMat)) { implicit builder => (leftSink, rightSink) =>
      import akka.stream.scaladsl.GraphDSL.Implicits._
      val partition = builder.add(Partition[FormData.BodyPart](2, element => if (element.getName == "request") 0 else 1))
      source ~> partition.in
      partition.out(0).map(_.getEntity) ~> leftSink
      partition.out(1).map(_.getEntity) ~> rightSink
      ClosedShape
    }
  }
}

第二步:调用分区方法,直接获取目标结果

不需要提前把Flow和Sink绑定,直接传入对应功能的Sink运行流即可:

// 直接定义两类Sink,对应你要的输出类型
val requestSink: Sink[BodyPartEntity, Future[Seq[BodyPartEntity]]] = Sink.seq[BodyPartEntity]
val dataSink: Sink[BodyPartEntity, Publisher[BodyPartEntity]] = Sink.asPublisher[BodyPartEntity](fanout = false)

// 运行你自己的FormData Source,拿到两个目标结果
val (requestSeqFuture, dataPublisher) = RunnableGraph.fromGraph(yourFormDataSource.partition(requestSink, dataSink)).run()

// 等待Future完成即可获取Seq[BodyPartEntity]
import scala.concurrent.ExecutionContext.Implicits.global
requestSeqFuture.foreach { requestSeq =>
  // 此处直接使用得到的requestSeq
}

核心逻辑说明

  • Sink.seq的物化值本身就是Future[Seq[Element]],流运行完成后Future会触发完成,拿到你要的序列
  • Sink.asPublisher的物化值直接就是Publisher[Element],流启动后即可直接使用该Publisher

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.23 22:24:03