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

如何在Akka Streams中将Flow1的物化值传递给Flow2使用?

How to Pass the Materialized Value of the First Flow to the Second Flow in Akka Streams

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() or materialize()).
  • 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 06:58:25