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

