Scala Spark流作业调用.toDF转换DataFrame空指针异常排查
问题根因
- 空指针核心触发原因:你将
eventList.toDF逻辑写在了DStream的map算子内部,DStream的算子计算逻辑运行在集群Executor节点,而你在Driver端experiment方法中导入的spark.implicits._对应的隐式编码器、SparkSession实例属于Driver端对象,不会被正确序列化到Executor。Executor执行到toDF调用时,localSeqToDatasetHolder内部依赖的SparkSession引用为null,直接抛出空指针。 - 附带逻辑错误:你使用
map(event => List.fill(1)(event)).reduce((a, b) => a ++ b)做聚合是错误实现,DStream的reduce是跨批次全局状态聚合,会将所有历史批次的数据一直累加在内存中,既不会按15秒窗口触发下游计算,还会最终导致内存溢出。 - 代码缺失问题:你反序列化JSON时使用的
Event样例类未定义,和你声明的SaveableEvent类不匹配。
修复方案
- 首先将所有样例类定义移到顶级作用域(不要定义在方法内部),避免隐式推导、序列化异常,补上缺失的
Event类:
import java.time.LocalDateTime final case class Coords(latitude: Double, longitude: Double) final case class Person(name: String, peacescore: Double) final case class Event( peacewatcherID: Int, timestamp: LocalDateTime, location: Coords, words: List[String], persons: List[Person] ) final case class SaveableEvent( peacewatcherID : Int, timestamp: String, location: Coords, words: List[String], persons: List[Person] )
- 替换原有流处理逻辑,使用
foreachRDD按批次处理数据,不要在Executor侧执行DataFrame转换、S3写入逻辑,同时删掉错误的全局reduce聚合逻辑,不要手动将分布式数据收集到Driver端拼成List:
stream .flatMap(record => { implicit val personFormat = Json.format[Person] implicit val coordsFormat = Json.format[Coords] implicit val eventFormat = Json.format[Event] val json = Json.parse(record.value()) eventFormat.reads(json).asOpt }) .map(event => { val serializedTime = event.timestamp.format(DateTimeFormatter.ISO_DATE_TIME) SaveableEvent(event.peacewatcherID, serializedTime, event.location, event.words, event.persons) }) .foreachRDD { rdd => // 每个批次在Driver端获取SparkSession,导入隐式转换 val spark = SparkSession.builder.config(rdd.sparkContext.getConf).getOrCreate() import spark.implicits._ // 空批次直接跳过 if (!rdd.isEmpty()) { // 直接将分布式RDD转为DataFrame,不需要collect到本地内存 val df = rdd.toDF() // 建议按批次时间设置子路径,避免覆盖历史数据 val batchPath = s"s3a://path/batch_time=${LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyyMMddHHmmss"))}" df.write.mode(SaveMode.Overwrite).parquet(batchPath) } }
额外优化建议
- 禁止手动将RDD数据
toList收集到Driver端再转DataFrame,数据量稍大就会触发Driver端OOM,直接基于RDD转DF是分布式执行,性能和稳定性更高。 - 如果使用Spark 2.3及以上版本,更建议替换为Structured Streaming实现Kafka消费、窗口聚合、S3写入逻辑,原生支持Encoder自动推导、Exactly-Once语义,不需要手动维护DStream的批次状态,代码量更少也更稳定。
内容的提问来源于stack exchange,提问作者SomeDev
相关产品推荐
相关产品推荐

