Akka广播流遇~>方法重载错误,求解决方案
问题原因及解决办法
你遇到的overloaded method ~> with alternatives错误,本质是类型不匹配导致的,具体有以下几个问题:
- 源输出是
ByteString类型,但你创建的Broadcast[String]接收的是String类型,二者无法直接连接 Flow[String]的输出类型和Zip[Type1, Type2]的输入类型不匹配(Type1和Type2需要是具体的、与Flow输出对应的类型)
具体修正步骤:
第一步:将ByteString转换为String
在源和广播器之间添加转换操作,把ByteString解码为String:byteStringSource.map(_.utf8String) ~> broadcast第二步:对齐Flow与Zip的类型
确保incrementer和multiplier的输出类型分别匹配Type1和Type2。比如假设Type1是String,Type2是String,可以直接使用;如果是其他类型(比如Int、自定义类),需要在Flow的map中做转换。第三步:确保Type1和Type2已定义
如果Type1和Type2是自定义类型,要提前声明;如果是内置类型,直接替换即可。
修正后的完整代码示例
假设Type1和Type2都是String(如果是其他类型,调整Flow中的map逻辑即可):
import akka.stream.scaladsl.{Broadcast, Flow, RunnableGraph, Sink, Source, Zip} import akka.stream.{ClosedShape, GraphDSL} import akka.util.ByteString // 假设Type1和Type2是String类型,若为自定义类型需提前定义 type Type1 = String type Type2 = String val byteStringSource: Source[ByteString, Any] = Source.fromIterator(() => (1 to 10).map(i => ByteString(s"Element $i")).iterator) val incrementer = Flow[String].map { x => // 这里可以添加Type1的转换逻辑,比如x => x.toInt(如果Type1是Int) x } val multiplier = Flow[String].map { x => // 同理添加Type2的转换逻辑 x } val output = Sink.foreach[(Type1, Type2)] { n1 => println(s"First obj is ${n1._1.toString} & second obj is ${n1._2.toString}") } val graph = RunnableGraph.fromGraph( GraphDSL.create() { implicit builder: GraphDSL.Builder[NotUsed] => import GraphDSL.Implicits._ val broadcast = builder.add(Broadcast[String](2)) val zip = builder.add(Zip[Type1, Type2]) // 关键:先将ByteString转成String再传给广播器 byteStringSource.map(_.utf8String) ~> broadcast broadcast.out(0) ~> incrementer ~> zip.in0 broadcast.out(1) ~> multiplier ~> zip.in1 zip.out ~> output ClosedShape } ) graph.run()
如果Type1和Type2是其他类型,比如Type1 = Int、Type2 = Double,只需修改Flow的map逻辑:
type Type1 = Int type Type2 = Double val incrementer = Flow[String].map { x => x.split(" ")(1).toInt + 1 // 提取数字并加1,输出Int } val multiplier = Flow[String].map { x => x.split(" ")(1).toDouble * 2 // 提取数字并乘2,输出Double }
内容的提问来源于stack exchange,提问作者Sushant Somani
相关产品推荐
相关产品推荐

