Scala2.12+Spark3.2中partitionIterator.nonEmpty报错求解决方案
问题原因
在Scala 2.12 + Spark 3.2环境下,foreachPartition接收的迭代器是Java Iterator,而非Scala原生Iterator。Scala 2.11中存在隐式转换可自动将Java Iterator转为Scala Iterator,因此能调用nonEmpty方法;但Scala 2.12中该隐式转换不再自动生效,直接调用nonEmpty就会触发报错。
解决方案
有两种可行的修复方式:
方式一:显式转换为Scala Iterator
导入Scala的Java迭代器转换工具,将Java Iterator转为Scala Iterator,即可正常使用nonEmpty方法:
import scala.collection.JavaConverters._ val ds = spark.createDataset(List("abc","efg","xyz")) ds.foreachPartition(partitionIterator => { val scalaIterator = partitionIterator.asScala if(scalaIterator.nonEmpty){ print(scalaIterator.getClass) } })
方式二:使用Java Iterator原生方法判断非空
直接用Java Iterator的hasNext()方法替代nonEmpty,无需额外转换:
val ds = spark.createDataset(List("abc","efg","xyz")) ds.foreachPartition(partitionIterator => { if(partitionIterator.hasNext()){ print(partitionIterator.getClass) } })
注意:
hasNext()不会消耗迭代器元素,和nonEmpty的行为完全一致,适合用来判断分区是否为空。
内容的提问来源于stack exchange,提问作者Surya
相关产品推荐
相关产品推荐

