Spark Dataset调用mapPartitions后show()显示空表问题排查
这种Dataset返回空表但RDD正常的情况,大概率是编码器(Encoder)不匹配导致的——毕竟RDD依赖Java/Kryo序列化,对数据格式要求宽松;而Dataset依赖Spark专用编码器做优化,一旦编码器和数据结构不兼容,就可能出现数据无法正确序列化/反序列化,最终show()看不到结果的情况。
咱们来拆解你代码里的几个关键问题:
1. 编码器选择错误
你用Encoders.bean(classOf[(String,MyState)])来编码Scala元组,但Encoders.bean()是专为JavaBean设计的,对Scala Tuple的支持非常差,尤其是和自定义case class结合时。而且你的MyState是Scala case class,应该用Spark为product类型提供的Encoders.product,而非bean编码器。
另外,MyState里使用了java.util.ArrayList,在Spark 2.2.0这个旧版本中,默认的case class编码器对Java集合的序列化支持并不完善,这也可能导致数据丢失。
2. 迭代器消耗问题(排除项)
虽然你提到调试时能看到返回结果,但要确认:executeStrategy里iter.toList会一次性消耗整个迭代器——不过RDD和Dataset的mapPartitions处理逻辑在这一点上是一致的,且你说RDD正常,所以这个不是核心问题。
具体修复步骤
步骤1:调整MyState的集合类型(推荐)
把Java集合换成Scala原生集合,让Spark默认编码器能更好地处理:
// 将java.util.ArrayList替换为Scala原生集合 case class MyState(code: List[Any], evaluation: List[Double])
如果必须保留Java集合,需要为MyState显式指定支持Java集合的编码器,但用Scala集合会更省心。
步骤2:修正编码器的使用
去掉自定义的enc1,让Spark自动推导编码器——因为MyState是case class,Spark会自动生成Encoders.product[MyState],元组(String, MyState)的编码器也会被自动推导:
def run(sparkSession: SparkSession, n: Int): Unit = { import sparkSession.implicits._ var data0 = conf.dataBase.repartition(n).persist(StorageLevel.MEMORY_AND_DISK_SER) // 移除自定义编码器,Spark会自动处理类型推导 val data1 = data0.mapPartitions(executeStrategy).toDF() data1.show(3) }
如果需要显式指定,应该为MyState单独指定编码器,而非元组:
implicit val myStateEncoder: Encoder[MyState] = Encoders.product[MyState] val data1 = data0.mapPartitions(executeStrategy).toDF()
步骤3:验证返回逻辑
可以在executeStrategy末尾加打印日志,确认每个分区的输出数量:
// 在executeStrategy最后一行前添加 println(s"Partition returning ${result.size} elements") result.map(x => (x.getCode.toString, MyState(x.getCode,x.getEvaluation))).iterator
额外提示:Spark版本问题
Spark 2.2.0确实比较旧,存在一些编码器相关的bug。如果有条件,升级到Spark 2.4.x(依然支持Scala 2.11)会减少这类序列化问题。
为什么RDD正常?
RDD使用的是Java序列化(默认)或Kryo序列化,只要你的MyState和元组能被Java序列化,就能正常输出。但Dataset依赖的是Spark专用编码器,它需要精确匹配数据结构,一旦编码器不兼容,数据会因序列化失败“消失”(Spark可能不会抛出明显错误,仅返回空结果)。
内容的提问来源于stack exchange,提问作者Alex100

