Spark中调用unpersist()触发序列化异常问题咨询
首先得明确为啥会出现这个错误:你在Executor的转换函数call()方法里调用广播变量的unpersist()操作,这个方法本质上会触发Spark内部的消息传递(就是报错里的RemoveBroadcast消息),而这个过程中涉及到的scala.concurrent.impl.Promise$DefaultPromise对象是不可序列化的。当Spark尝试把相关对象序列化传递时,就抛出了NotSerializableException。
而且从Spark的设计逻辑来说,广播变量的生命周期管理(比如持久化、清理)本来就该由Driver端来控制,Executor的职责只是读取和使用广播变量,不应该去修改它的持久化状态。
给你几个可行的解决办法:
最推荐的方式:移到Driver端执行unpersist
等所有依赖这个广播变量的Spark任务都执行完成后,在Driver代码里调用broadcastVar.unpersist()(或者unpersist(true)强制清理)。比如:// 执行所有使用广播变量的任务 val result = rdd.map(...).collect() // 任务完成后在Driver端清理 broadcastHashMap.unpersist()如果确实需要在Executor端触发清理(不推荐,除非特殊场景)
可以通过SparkContext的异步操作间接实现,但绝对不要在任务的call方法里直接调用unpersist。比如可以在Executor里向Driver发送信号,让Driver去执行清理操作,但这种方式会增加代码复杂度,一般不建议这么做。排查代码中的序列化隐患
检查下你的代码里有没有不小心把Driver端的非序列化对象(比如Promise、未序列化的自定义类)传递到Executor的任务逻辑里,这也可能间接引发这类序列化异常。
内容的提问来源于stack exchange,提问作者Rushi Pradhan

