【问题标题】:Feeding the output of a Flow to a Broadcast in Akka Streams Graph将流的输出馈送到 Akka 流图中的广播
【发布时间】:2017-08-31 06:20:18
【问题描述】:

我正在尝试编写 Akka 流图。我写的代码是

val graph = RunnableGraph.fromGraph(GraphDSL.create(sink1, sink2)((_, _)) { implicit builder =>
   (sink1, sink2) =>
      import GraphDSL.Implicits._
      val bcast = builder.add(Broadcast[Row](2))
      val flow = source ~> flow1 ~> flow2
      flow.out ~> bcast.in
      bcast.out(0) ~> sink1
      bcast.out(1) ~> flow3 ~> flow4 ~> sink2
      ClosedShape
})

val (f1, f2) = graph.run()
val consolidated = Future.sequence(List(f1, f2))
Await.result(consolidated, Duration.Inf)

此代码无法编译,因为我无法将流出的流连接到 bcast 的流入。

我可以将源的输出连接到 bcast 的输入,但我不能这样做,因为两个分支之间的某些部分是共同的。所以我必须在flow2之后才在图中创建分支

另外...我不确定我是否正确编写了图表,因为它返回了两个完成的未来,我需要使用序列手动将它们组合成一个未来。

【问题讨论】:

    标签: scala akka-stream


    【解决方案1】:

    您不能分两步连接您的图表,因为~> 组合器不会给您返回一个流程。它实际上是一个有状态的声明式操作。

    这里更好的方法是一次性连接您的图表,例如

      source ~> flow1 ~> flow2 ~> bcast
                                  bcast          ~>          sink1
                                  bcast ~> flow3 ~> flow4 ~> sink2
    

    或者,您也可以通过向构建器添加一个阶段(并检索其形状)来拆分声明,例如

      val flow2s = builder.add(flow2)
    
      source ~> flow1 ~> flow2s.in
      flow2s.out ~> bcast
                    bcast          ~>          sink1
                    bcast ~> flow3 ~> flow4 ~> sink2
    

    关于物化Futures,您需要选择有意义的东西作为整个图形的物化值。如果您只需要物化 Futures 的 2 个 Sinks 之一,则只需将那一个传递给 GraphDSL.create 方法。 否则,如果您对Futures 都感兴趣,那么将sequencezip 放在一起非常有意义。

    【讨论】:

      猜你喜欢
      • 2019-06-25
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多