【发布时间】: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