集群模式下无法访问RDD,执行foreach时抛出空指针异常
集群模式下Spark RDD foreach空指针问题解决
问题场景
用以下代码从Seq创建RDD:
val testRDD = sc.parallelize(list.toSeq, size);
执行遍历操作时,集群模式抛出空指针异常,客户端模式运行正常:
testRDD.foreach(row => { logger.info("Row index " + row.index.toString() ); })
补充信息:
testRDD.count()和testRDD.partitions.size返回结果正常- 执行
collect()后再调用foreach可正常运行,但因需分布式处理,不想使用collect()
问题原因
- 日志实例无法跨节点序列化:客户端模式下,Driver和Executor在同一进程,闭包里直接引用的Driver端
logger可以正常调用;但集群模式下,Executor是独立进程,logger属于Driver端的非序列化对象,无法被传递到Executor,调用时触发空指针。 - 自定义Row对象序列化问题:如果
row是自定义类,未实现Serializable接口的话,集群模式下数据序列化传输时可能导致对象属性初始化不全,访问row.index时抛出空指针。
解决办法
1. 在Executor端本地初始化日志实例
不要直接引用Driver端的logger,而是在foreach闭包内创建Executor本地的日志实例,示例代码:
import org.slf4j.LoggerFactory testRDD.foreach(row => { val localLogger = LoggerFactory.getLogger(getClass) localLogger.info("Row index " + row.index.toString()) })
2. 确保自定义Row对象可序列化
如果row是自定义类,必须实现Serializable接口,示例:
class MyRow(val index: Int) extends Serializable { // 类的其他逻辑 }
3. 避免在闭包中引用Driver端非序列化资源
所有在foreach闭包中用到的资源,要么在Executor本地初始化,要么确保资源对象实现了Serializable接口,禁止直接引用Driver端的非序列化实例。
内容的提问来源于stack exchange,提问作者Faizal
相关产品推荐
相关产品推荐

