如何在Akka Streams中将Flow1的物化值传递给Flow2使用?
Hey there! I totally get the confusion around materialized values—they're one of the trickier concepts in Akka Streams, especially when working with custom GraphStage implementations. Let's break down exactly how to pass a materialized value from your first flow to the second, with practical examples tailored to your use case.
First, a Quick Refresher
A materialized value is an object generated when the stream starts running (not when you build the flow). It’s usually a handle to interact with the stream (like an ActorRef for sending control messages, a Promise for tracking completion, or a counter for metrics). The challenge here is that you can’t access this value while building your flow topology—you have to wait until the stream is materialized.
Two Common Solutions
1. Embed the Materialized Value in Stream Elements
If your materialized value is immutable and can be attached to each stream element, this is the simplest approach. You’ll modify the first flow to emit tuples of (OriginalElement, MaterializedValue), then the second flow can unpack and use the value.
Example with a Custom GraphStage
import akka.stream._ import akka.stream.stage._ // First Flow: Generates a unique ID as its materialized value, emits (element, ID) class FirstFlow extends GraphStageWithMaterializedValue[FlowShape[String, (String, String)], String] { val in = Inlet[String]("FirstFlow.in") val out = Outlet[(String, String)]("FirstFlow.out") override val shape = FlowShape(in, out) override def createLogicAndMaterializedValue(attrs: Attributes): (GraphStageLogic, String) = { // Generate a unique ID as the materialized value val streamId = java.util.UUID.randomUUID().toString val logic = new GraphStageLogic(shape) { setHandler(in, new InHandler { override def onPush(): Unit = { val elem = grab(in) // Emit element + materialized value together push(out, (elem, streamId)) } }) setHandler(out, new OutHandler { override def onPull(): Unit = pull(in) }) } (logic, streamId) } } // Second Flow: Unpacks the tuple and uses the materialized value val secondFlow = Flow[(String, String)].map { case (elem, streamId) => s"Processed element '$elem' using stream ID: $streamId" } // Combine and run the streams val firstFlow = Flow.fromGraph(new FirstFlow()) val combinedFlow = firstFlow.via(secondFlow) combinedFlow.runWith(Source(List("foo", "bar")), Sink.foreach(println))
This works well for values that don’t change during stream execution, but it does alter your element type (you’ll need to handle tuples throughout the downstream flow).
2. Pass the Materialized Value Post-Materialization
If your materialized value is a mutable/interactive object (like an ActorRef or Promise), you can grab it after the stream starts, then send it to the second flow’s logic. This requires the second flow to support receiving external messages (via AsyncCallback in a GraphStage).
Example with ActorRef as Materialized Value
import akka.actor.ActorRef import akka.stream.scaladsl._ // First Flow: Materialized value is an ActorRef for control messages class FirstFlow extends GraphStageWithMaterializedValue[FlowShape[String, String], ActorRef] { val in = Inlet[String]("FirstFlow.in") val out = Outlet[String]("FirstFlow.out") override val shape = FlowShape(in, out) override def createLogicAndMaterializedValue(attrs: Attributes): (GraphStageLogic, ActorRef) = { // Create an AsyncCallback to handle control messages val controlCallback = createAsyncCallback[String] { msg => println(s"FirstFlow received control: $msg") } val logic = new GraphStageLogic(shape) { setHandler(in, new InHandler { override def onPush(): Unit = push(out, s"FirstFlow: ${grab(in)}") }) setHandler(out, new OutHandler { override def onPull(): Unit = pull(in) }) } // Return the ActorRef of the callback as the materialized value (logic, controlCallback.actorRef) } } // Second Flow: Can receive the ActorRef from the first flow class SecondFlow extends GraphStage[FlowShape[String, String]] { val in = Inlet[String]("SecondFlow.in") val out = Outlet[String]("SecondFlow.out") override val shape = FlowShape(in, out) // Callback to receive the first flow's ActorRef private val setControlActor = createAsyncCallback[ActorRef] { ref => controlActor = Some(ref) } // Expose an ActorRef to send the materialized value to this flow def setControlActorRef: ActorRef = setControlActor.actorRef private var controlActor: Option[ActorRef] = None override def createLogic(attrs: Attributes): GraphStageLogic = new GraphStageLogic(shape) { setHandler(in, new InHandler { override def onPush(): Unit = { val elem = grab(in) // Use the first flow's ActorRef if available controlActor.foreach(_ ! s"SecondFlow processing: $elem") push(out, s"SecondFlow: $elem") } }) setHandler(out, new OutHandler { override def onPull(): Unit = pull(in) }) } } // Build and run the stream implicit val system = ActorSystem("MatValueDemo") implicit val ec = system.dispatcher val firstFlow = Flow.fromGraph(new FirstFlow()) val secondFlow = new SecondFlow() // Build the topology val graph = RunnableGraph.fromGraph(GraphDSL.create(firstFlow) { implicit builder => flow1 => import GraphDSL.Implicits._ val flow2 = builder.add(secondFlow) flow1 ~> flow2 ClosedShape }) // Materialize the stream to get the first flow's ActorRef val firstFlowMatVal = graph.run() // Pass the materialized value to the second flow secondFlow.setControlActorRef ! firstFlowMatVal
Here, after starting the stream, we take the ActorRef from the first flow and send it to the second flow’s control endpoint. The second flow can then use this reference to interact with the first flow during runtime.
Key Takeaways
- Materialized values are only available after the stream is materialized (when you call
run()ormaterialize()). - Choose embedding if your value is static and can be attached to elements.
- Choose post-materialization passing if your value is interactive (like an
ActorRef) and needs to be sent to the downstream flow’s logic.
内容的提问来源于stack exchange,提问作者Nikhil

