调用RDD.mapValues后使用reduceByKey编译失败,调换顺序则正常
mapValues后调用reduceByKey编译失败的问题 我太懂这个坑了——这其实是Scala类型推断和Spark隐式转换机制共同作用下的典型问题,我之前也踩过类似的雷。先看你给出的复现代码:
def test[E]() = new SparkContext().textFile("").keyBy(_ ⇒ 0L).mapValues(_.asInstanceOf[E]).reduceByKey((x, _) ⇒ x)
编译报错提示reduceByKey不是当前RDD的成员,核心原因是Scala编译器无法正确推断泛型E的类型,导致Spark的隐式转换无法生效。
Spark的reduceByKey这类KV操作,本质是依赖PairRDDFunctions这个隐式转换类——当你的RDD是RDD[(K, V)]类型时,编译器会自动把它转换成PairRDDFunctions,从而让你调用这些KV专属方法。但在你的代码里,mapValues(_.asInstanceOf[E])这一步,因为E是一个未绑定的泛型参数,编译器没办法确定转换后的RDD具体类型是RDD[(Long, E)],自然找不到对应的隐式转换,也就识别不出reduceByKey方法。
而反过来先调用reduceByKey再mapValues时,reduceByKey处理的是明确的RDD[(Long, String)]类型,转换后的RDD类型也很清晰,后续mapValues的类型推断就没有障碍,所以编译正常。
几个可行的解决方案:
1. 显式指定mapValues的泛型类型
直接帮编译器明确mapValues的输出值类型,让它知道转换后的RDD是RDD[(Long, E)]:
def test[E]() = new SparkContext().textFile("").keyBy(_ ⇒ 0L).mapValues[E](_.asInstanceOf[E]).reduceByKey((x, _) ⇒ x)
2. 拆分链式调用,显式声明中间RDD类型
把一步链式调用拆成多步,给每个中间RDD指定明确的类型,消除编译器的推断歧义:
def test[E]() = { val sc = new SparkContext() val rawPairRDD: RDD[(Long, String)] = sc.textFile("").keyBy(_ ⇒ 0L) val typedPairRDD: RDD[(Long, E)] = rawPairRDD.mapValues(_.asInstanceOf[E]) typedPairRDD.reduceByKey((x, _) ⇒ x) }
3. 用map替代mapValues(类型推断更友好)
有时候map的类型推断逻辑比mapValues更直接,换写法也能解决问题:
def test[E]() = new SparkContext().textFile("").keyBy(_ ⇒ 0L).map{ case (k, v) => (k, v.asInstanceOf[E]) }.reduceByKey((x, _) ⇒ x)
这些方法本质都是帮编译器明确泛型类型,让Spark的隐式转换能正常触发,从而调用到reduceByKey方法。
内容的提问来源于stack exchange,提问作者Reinstate Monica

