【发布时间】:2021-01-06 05:19:07
【问题描述】:
我正在尝试 Akka Stream API,但我不知道为什么它会抛出 java.lang.IllegalArgumentException。
val graph = RunnableGraph.fromGraph(
GraphDSL.create(source, sink)
((source, sink) => Seq(source, sink)) {
implicit b => (source, sink) =>
Import akka.stream.scaladsl.GraphDSL.Implicits._
val partition = b.add(Partition[(KinesisRecord)](2, flow => {
1
}))
source ~> partition.in
partition.out(0) ~> sink
partition.out(1) ~> sink
ClosedShape
})
这是当前代码。错误如下
[info] - should consume *** FAILED ***
[info] java.lang.IllegalArgumentException: [Map.in] is already connected
[info] at akka.stream.scaladsl.GraphDSL$Builder.addEdge(Graph.scala:1567)
[info] at akka.stream.scaladsl.GraphDSL$Implicits$CombinerBase.$tilde$greater(Graph.scala:1730)
[info] at akka.stream.scaladsl.GraphDSL$Implicits$CombinerBase.$tilde$greater$(Graph.scala:1729)
[info] at akka.stream.scaladsl.GraphDSL$Implicits$PortOpsImpl.$tilde$greater(Graph.scala:1784)
[info] at akka.stream.scaladsl.GraphApply.create(GraphApply.scala:46)
[info] at akka.stream.scaladsl.GraphApply.create$(GraphApply.scala:41)
[info] at akka.stream.scaladsl.GraphDSL$.create(Graph.scala:1529)
我使用 kinesisRecord 作为源的目标。 但是,在这段代码中,如果我将 outputPorts 更改为 1 并删除
partition.out(1) ~> sink
这一行,行得通。
我不知道是我遗漏了什么还是只是一个错误。
【问题讨论】:
-
分区的输出需要
Merge[XXX](2)形状来合并2 个输出,而不是只有一个输入的Sink。这就是为什么如果您删除您说它有效的代码。在Merge之后,您连接到您的Sink。 -
完全正确。感谢您的尖锐评论。
-
什么是
Source,什么是Sink?
标签: scala akka akka-stream