求基于Apache Beam Java API的Scala Combine.perKey可用示例
在Scala中使用Apache Beam Java API的Combine.perKey解决示例
我明白你在Scala里用Apache Beam Java API的Combine.perKey时遇到了类型互操作性的头疼问题——Scala和Java的泛型规则差异确实容易在这里踩坑。咱们直接来看怎么修复你的代码,让它能正常运行:
问题根源分析
你的原代码报错主要有两个原因:
- 类型参数数量错误:Combine.perKey的Java签名只需要指定输入值类型和输出值类型,不需要把键的类型也传进去,你传了三个类型参数导致报错“too many type arguments”。
- 类型推断失效:移除类型参数后,Scala无法正确推断出Combine.perKey的OutputT类型,这通常是因为自定义的
SerializableFunction和Beam的Java类型衔接有问题。
修复后的可运行代码
方案1:使用自定义的SumLongs合并器
首先确保你导入了正确的Beam包,然后调整自定义合并器的类型(注意要实现Java的SerializableFunction,并使用Java的Iterable):
import org.apache.beam.sdk.transforms.Combine import org.apache.beam.sdk.transforms.SerializableFunction import org.apache.beam.sdk.values.KV import org.apache.beam.sdk.values.PCollection // 自定义合并器:实现Java的SerializableFunction,接收Java Iterable class SumLongs extends SerializableFunction[java.lang.Iterable[Long], Long] { override def apply(input: java.lang.Iterable[Long]): Long = { var sum = 0L val iterator = input.iterator() while (iterator.hasNext) { sum += iterator.next() } sum } } // 你的PCollection初始化(这里示例用占位符,实际替换成你的数据源) val sales: PCollection[KV[(Int, Int), Long]] = ??? // 正确调用Combine.perKey:要么让Scala自动推断类型,要么显式指定输入/输出值类型 val summedSales = sales.apply(Combine.perKey(new SumLongs())) // 或者显式指定类型参数(可选,帮助Scala推断) // val summedSales = sales.apply(Combine.perKey[Long, Long](new SumLongs()))
方案2:使用Beam内置的Sum合并器(更简洁)
Beam已经提供了现成的数值合并器,完全不需要自己写SumLongs,这样能避免类型互操作问题:
import org.apache.beam.sdk.transforms.Combine import org.apache.beam.sdk.transforms.Sum import org.apache.beam.sdk.values.KV import org.apache.beam.sdk.values.PCollection val sales: PCollection[KV[(Int, Int), Long]] = ??? // 直接使用内置的Sum.ofLongs(),类型推断会自动完成 val summedSales = sales.apply(Combine.perKey(Sum.ofLongs()))
额外说明
- Scala和Java的
Iterable是不同的类型,Beam的Java API需要接收Java的Iterable,所以自定义合并器里要使用java.lang.Iterable而不是Scala的Iterable。 - 当调用Java泛型方法时,Scala的类型推断有时候需要一点帮助,显式指定必要的类型参数(比如Combine.perKey的输入/输出值类型)就能解决“Unspecified type parameters”的问题。
内容的提问来源于stack exchange,提问作者Alistair McIntyre
相关产品推荐
相关产品推荐

