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

如何在Akka Stream中消费带双Sink的Source并获取其一结果

How to Retrieve the Materialized Value from a Sink in a Broadcast Akka Stream Graph

Absolutely feasible! The problem with your current code is that you aren’t capturing the materialized value of addSink when constructing your graph—right now, run() returns NotUsed because the graph builder has no instruction to retain the sink’s output. Let’s fix this step by step:

Corrected Code

import akka.stream.scaladsl.{Broadcast, GraphDSL, RunnableGraph, Sink, Source}
import scala.concurrent.Future

val source = Source(1 to 20)
val addSink = Sink.fold[Int, Int](0)(_ + _)
val subtractSink = Sink.fold[Int, Int](0)(_ - _)

// Pass addSink to GraphDSL.create to capture its materialized value
val graph = GraphDSL.create(addSink) { implicit builder => addSink =>
  import GraphDSL.Implicits._
  val bcast = builder.add(Broadcast[Int](2))
  
  source ~> bcast.in
  bcast.out(0) ~> addSink
  bcast.out(1) ~> subtractSink // Fixed typo: subtrackSink → subtractSink
  
  ClosedShape
}

// Now run() returns the materialized value of addSink: Future[Int]
val result: Future[Int] = RunnableGraph.fromGraph(graph).run()

Key Changes Explained

  • Capture the Sink’s Output: By passing addSink as an argument to GraphDSL.create, you tell the graph builder to retain and return this sink’s materialized value (which is a Future[Int] for Sink.fold).
  • Typo Fix: Corrected subtrackSink to subtractSink to resolve compilation errors.

Bonus: Get Results from Both Sinks

If you ever need to retrieve outputs from both sinks, you can pass both to GraphDSL.create and return a tuple of their materialized values:

val graph = GraphDSL.create(addSink, subtractSink) { implicit builder => (addSink, subtractSink) =>
  import GraphDSL.Implicits._
  val bcast = builder.add(Broadcast[Int](2))
  
  source ~> bcast.in
  bcast.out(0) ~> addSink
  bcast.out(1) ~> subtractSink
  
  ClosedShape
}

val (addResult, subtractResult): (Future[Int], Future[Int]) = RunnableGraph.fromGraph(graph).run()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 08:51:43