Spark Scala报错:mapPartitionsWithIndex不是RDD[Int]的成员
问题分析与解决
嘿,我来帮你排查这个报错!你遇到的问题根源在于myfunc函数的类型声明错误,导致Spark无法找到匹配的mapPartitionsWithIndex方法。
错误原因
你的代码里,myfunc的第二个参数写的是iter: Iterator[(int)],这里有两个关键问题:
- Scala中的整数类型是大写的
Int,不是Java里的小写int; - 多了一对不必要的括号,正确的迭代器类型应该是
Iterator[Int]。
mapPartitionsWithIndex方法要求传入的函数签名是(Int, Iterator[T]) => Iterator[U],而你的RDD是RDD[Int],所以迭代器的类型必须是Iterator[Int]。类型不匹配的情况下,Spark就会报错说该方法不是RDD[Int]的成员。
修正后的代码
val z = sc.parallelize(List(1,2,3,4,5,6),2) // 打印带分区标签的内容 def myfunc(index: Int, iter: Iterator[Int] ): Iterator[String] = { iter.toList.map( x => "[Part Id: " + index + " ,val:" + x + "]").iterator } z.mapPartitionsWithIndex(myfunc).collect
验证效果
运行修正后的代码,你会得到类似这样的输出:
Array([Part Id: 0 ,val:1], [Part Id: 0 ,val:2], [Part Id: 0 ,val:3], [Part Id: 1 ,val:4], [Part Id: 1 ,val:5], [Part Id: 1 ,val:6])
内容的提问来源于stack exchange,提问作者M_Gandhi
相关产品推荐
相关产品推荐

