在实际项目中,序列化相关的报错非常常见,尤其是在 Spark Streaming 开发中。本文将深入解析 foreachRDD、foreachPartition 和 foreach 这三个方法的区别,帮助开发者快速定位并解决序列化问题。
这三个方法的核心差异体现在作用范围上:foreachRDD 作用于 DStream 中每个时间间隔的 RDD;foreachPartition 作用于每个时间间隔的 RDD 中的每个分区;而 foreach 则作用于每个时间间隔的 RDD 中的每个元素。虽然概念看似复杂,但关键在于它们执行位置的差异——这直接决定了序列化行为。
从执行层面来看,foreachRDD 运行在 driver 端,而 foreachPartition 和 foreach 则运行在 worker 端。这一点至关重要——如果在 worker 端错误地使用了 driver 端的对象,就会引发序列化异常。例如,使用 foreachRDD 向外部系统输出数据时,通常需要创建连接对象。如果像下面这样将连接创建在 driver 端,那么 foreach 在每个 worker 节点上执行时,节点上并不存在该连接对象,从而导致序列化错误或初始化失败。
dstream.foreachRDD { rdd =>
val connection = createNewConnection()
rdd.foreach { record =>
connection.send(record) // executed at the worker
}
}

正确的做法如下:
dstream.foreachRDD { rdd =>
rdd.foreachPartition { partitionOfRecords =>
val connection = createNewConnection()
partitionOfRecords.foreach(record => connection.send(record))
connection.close()
}
}
因此,driver 与 worker 之间的通信必须经过序列化。然而,并非所有对象都支持序列化。大多数序列化异常的场景,根源都在于 driver 端创建了不可序列化的对象,却试图在 worker 端使用它。掌握这个原则,排查序列化问题将变得清晰高效。
