使用Akka GraphDSL时Zip阶段的语法问题咨询
Akka GraphDSL Zip阶段端口连接问题
问题场景
尝试用Akka GraphDSL构建流图时,写出如下代码:
GraphDSL.create() { implicit builder => import GraphDSL.Implicits._ val in = Source(0 to 10) val fanOut = builder.add(Broadcast[Int](2)) val toString = builder.add(Flow[Int].map(_.toString)) val squared = builder.add(Flow[Int].map(x => x * x)) val zip = builder.add(ZipLatestWith((str: String, sqr: Int) => (str, sqr))) val out = Sink.ignore in ~> fanOut ~> toString ~> zip ~> out fanOut ~> squared ~> zip ClosedShape }
但在zip ~> out处报错:overloaded method ~> can't be applied to FanInShape2[String, Int, (String, Int)]。
改用以下写法后可正常运行:
in ~> fanOut ~> toString ~> zip.in0; zip.out ~> out fanOut ~> squared ~> zip.in1
疑问:见过部分教程无需指定.in和.out端口即可定义分支,想咨询这是GraphDSL对Zip阶段存在限制,还是写法存在错误?
原因解析
这不是写法错误,而是GraphDSL针对不同节点类型的端口隐式处理规则导致的:
单一输入/输出节点的隐式处理:对于Source、Flow、Sink,或者Broadcast这类FanOutShapeN节点,GraphDSL的隐式转换会自动处理端口选择:
- 比如
fanOut ~> toString会自动绑定Broadcast的第一个输出端口(out0)到Flow的输入;第二个fanOut ~> squared则自动绑定out1端口 - Flow链式连接
a ~> b ~> c会自动用前一个节点的输出端口连接后一个的输入端口
- 比如
FanInShape节点的显式输出要求:像ZipLatestWith、Zip这类多输入合并的节点属于FanInShape类型,它们的Shape包含多个输入端口和一个输出端口。GraphDSL不会自动将整个节点实例视为输出端口,因此必须显式指定
.out来获取输出端口,才能继续连接到Sink。
至于输入端口的隐式绑定(比如toString ~> zip自动连到in0,squared ~> zip自动连到in1),是GraphDSL的隐式规则会按顺序分配FanInShape的输入端口,这也是部分教程无需指定.in0/.in1的原因。
内容的提问来源于stack exchange,提问作者Uko
相关产品推荐
相关产品推荐

