启用Spark Streaming Checkpoint后遭遇java.io.NotSerializableException问题
咱先搞明白问题根源:当你开启Checkpoint后,Spark需要把整个流作业的状态(包括算子中用到的所有对象)序列化后持久化到存储介质里。但com.fasterxml.jackson.module.paranamer.shaded.CachingParanamer这个类根本没实现Serializable接口,而你的依赖类里刚好持有这个对象的引用,所以序列化直接失败——没开Checkpoint的时候,Spark不需要序列化全量状态,自然就不会触发这个问题。
给你几个靠谱的解决思路,按优先级排序:
1. 给非序列化成员加transient标记
如果这个CachingParanamer是你自己代码里依赖类的成员变量,直接给它加上transient关键字,告诉序列化机制跳过这个字段:
// 原代码 private CachingParanamer paranamer; // 修改后 private transient CachingParanamer paranamer;
要是Checkpoint恢复后这个变量需要重新初始化,记得在类里加个readObject方法手动恢复实例:
private void readObject(ObjectInputStream in) throws IOException, ClassNotFoundException { in.defaultReadObject(); // 重新创建paranamer实例 paranamer = new CachingParanamer(); }
2. 隔离第三方库的非序列化对象
如果这个依赖类是第三方库的,没法修改源码,那可以用ThreadLocal来封装这个非序列化对象,让它只存在于线程上下文里,Checkpoint序列化时不会处理ThreadLocal内的内容:
private ThreadLocal<CachingParanamer> paranamerHolder = ThreadLocal.withInitial(() -> new CachingParanamer()); // 使用时从ThreadLocal中获取 public void doBusinessLogic() { CachingParanamer paranamer = paranamerHolder.get(); // ...你的业务代码 }
3. 切换到Kryo序列化器
Spark默认的Java序列化器对非Serializable类支持很差,换成Kryo序列化器试试,它对多数类都能直接序列化,且效率更高。在Spark配置里这么修改:
val sparkConf = new SparkConf() .setAppName("YourStreamingApp") .set("spark.serializer", "org.apache.spark.serializer.KryoSerializer") // 将CachingParanamer注册到Kryo中 .registerKryoClasses(Array(classOf[com.fasterxml.jackson.module.paranamer.shaded.CachingParanamer]))
如果Kryo默认处理不了,还能给这个类写自定义Kryo序列化器,不过一般情况下注册后就能解决问题。
4. 检查Checkpoint的必要性
如果你的流作业是无状态的(没有使用updateStateByKey、窗口聚合这类需要保存状态的算子),可以考虑关闭Checkpoint;如果必须开启,确认是否可以只对必要状态做Checkpoint,缩小序列化范围。
内容的提问来源于stack exchange,提问作者fledgling

