升级Spark-Snowflake连接器后出现NotSerializableException异常求助
问题成因及解决方法
成因分析
- 版本兼容性不匹配:AWS EMR 6.6.0内置Spark 3.2.x,升级后的
spark-snowflake_2.12:2.11.0-spark_3.2连接器内部逻辑引入了对Spark内部非序列化类org.apache.spark.storage.StorageStatus的引用。该类未实现Serializable接口,Spark跨节点RPC通信(如BlockManagerMasterEndpoint交互)时需序列化此对象,直接触发NotSerializableException。 - 连接器逻辑变更:新版本连接器新增的存储状态监控或数据交互逻辑中,意外持有
StorageStatus实例引用,导致该对象被纳入序列化范围。应用虽能通过降级序列化方式(从Kryo fallback到Java序列化)继续运行,但Java序列化性能远低于Kryo,直接导致运行速度大幅下降。
解决方法
- 降级连接器版本:切换至与EMR 6.6.0 Spark 3.2兼容的稳定版本,例如
spark-snowflake_2.12:2.10.0-spark_3.2。该版本未引入导致序列化问题的逻辑,可直接规避此异常。 - 排查代码逻辑:检查应用中是否存在将Spark存储状态(如缓存DataFrame的状态)与Snowflake读写操作结合的代码,确保
StorageStatus相关对象仅在Driver端处理,不传递至Executor节点。例如,避免在UDF或分布式操作中引用存储状态信息。 - 强制启用Kryo序列化并优化配置:在Spark配置中指定使用Kryo序列化器,同时注册必要的类,减少Java序列化的 fallback 概率:
注意:由于spark.serializer=org.apache.spark.serializer.KryoSerializer spark.kryo.registerClasses=org.apache.spark.storage.StorageStatus spark.kryo.unsafe=trueStorageStatus是Spark内部类,此方法仅能缓解问题,无法保证彻底解决,需结合实际测试验证。 - 恢复EMR Spark默认配置:若自定义了BlockManager或RPC相关的Spark配置,尝试恢复为EMR默认配置,排查是否存在与新版本连接器的冲突项。
内容的提问来源于stack exchange,提问作者user19627118
相关产品推荐
相关产品推荐

