Spark连接MongoDB报ClassCastException:List序列化代理类型转换错误求助
Spark读取MongoDB时ClassCastException问题排查
问题场景
本地IDEA运行Spark Driver,Spark Master和Worker节点部署在本地虚拟机的Docker容器中,尝试读取MongoDB数据时抛出ClassCastException,初步怀疑依赖问题但无法定位具体冲突,分析Maven依赖树未发现明显异常。
异常信息
23/04/20 14:48:03 WARN TaskSetManager: Lost task 0.0 in stage 0.0 (TID 0) (172.19.0.6 executor 0): java.lang.ClassCastException: cannot assign instance of scala.collection.immutable.List$SerializationProxy to field org.apache.spark.sql.execution.datasources.v2.DataSourceRDDPartition.inputPartitions of type scala.collection.Seq in instance of org.apache.spark.sql.execution.datasources.v2.DataSourceRDDPartition at java.base/java.io.ObjectStreamClass$FieldReflector.setObjFieldValues(Unknown Source) at java.base/java.io.ObjectStreamClass$FieldReflector.checkObjectFieldValueTypes(Unknown Source) at java.base/java.io.ObjectStreamClass.checkObjFieldValueTypes(Unknown Source) at java.base/java.io.ObjectInputStream.defaultCheckFieldValues(Unknown Source) at java.base/java.io.ObjectInputStream.readSerialData(Unknown Source) at java.base/java.io.ObjectInputStream.readOrdinaryObject(Unknown Source) at java.base/java.io.ObjectInputStream.readObject0(Unknown Source) at java.base/java.io.ObjectInputStream.defaultReadFields(Unknown Source) at java.base/java.io.ObjectInputStream.readSerialData(Unknown Source) at java.base/java.io.ObjectInputStream.readOrdinaryObject(Unknown Source) at java.base/java.io.ObjectInputStream.readObject0(Unknown Source) at java.base/java.io.ObjectInputStream.readObject(Unknown Source) at java.base/java.io.ObjectInputStream.readObject(Unknown Source) at org.apache.spark.serializer.JavaDeserializationStream.readObject(JavaSerializer.scala:87) at org.apache.spark.serializer.JavaSerializerInstance.deserialize(JavaSerializer.scala:129) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:507) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source) at java.base/java.lang.Thread.run(Unknown Source)
已执行操作
已在Docker中Spark节点的/opt/spark/jars目录添加mongo-spark-connector_2.12:10.1.1依赖。
配置信息
SparkConf代码
public static SparkConf buildSparkConf() { return new SparkConf() .setMaster("spark://192.168.56.1:9013") .setAppName("test") .set("spark.local.ip","192.168.56.1") .set("spark.ui.port","9016") .set("spark.driver.bindAddress", "0.0.0.0") .set("spark.driver.port", "9360") .set("spark.blockManager.port", "9260") .set("spark.sql.shuffle.partitions", "30") .set("spark.sql.files.maxPartitionBytes", "1073741824") .set("spark.sql.broadcastTimeout", "36000") .set("spark.mongodb.input.uri", "mongodb://192.168.56.1/test.bigData") .set("spark.mongodb.write.connection.uri", "mongodb://192.168.56.1/test.bigData") .set("spark.mongodb.read.connection.uri", "mongodb://192.168.56.1/test.bigData") .set("spark.jars.packages", "org.mongodb.spark:mongo-spark-connector_2.12:10.1.1") .set("spark.driver.host","192.168.56.1"); } public void test() { SparkSession sparkSession = SparkSession.builder() .config(buildSparkConf()) .getOrCreate(); sparkSession.read() .format("mongodb") .load() .show(false); }
Maven依赖配置
<dependency> <groupId>org.mongodb</groupId> <artifactId>mongodb-driver-sync</artifactId> <version>4.8.2</version> </dependency> <dependency> <groupId>org.mongodb.spark</groupId> <artifactId>mongo-spark-connector_2.12</artifactId> <version>10.1.1</version> <exclusions> <exclusion> <artifactId>mongodb-driver-sync</artifactId> <groupId>org.mongodb</groupId> </exclusion> </exclusions> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.12</artifactId> <version>3.3.1</version> </dependency> <dependency> <!-- Spark dependency --> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.12</artifactId> <version>3.3.1</version> <scope>compile</scope> </dependency>
环境版本
- Scala版本:2.12.15
- Spark版本:3.3.1
- MongoDB版本:4.4.20
排查方向
- Scala版本一致性检查:确认Docker中Spark节点使用的Scala版本是否为2.12.x,与本地项目的2.12.15完全匹配。跨Scala版本编译的组件会导致序列化类结构不兼容,引发该类型转换异常。可进入Docker容器执行
scala -version或查看Spark安装包命名确认版本。 - 依赖部署方式冲突:本地代码中同时通过
spark.jars.packages指定Connector,又手动在Docker节点添加依赖,可能导致版本重复或不一致。建议只保留一种依赖部署方式:要么依赖集群自动拉取,要么确保本地与Docker中的Connector版本完全一致且无重复。 - Scala库依赖一致性:检查本地项目是否引入了多版本的
scala-library,执行mvn dependency:tree | grep scala-library确认仅存在2.12.15版本;同时检查Docker中Spark自带的scala-library版本,确保与本地一致,必要时替换Docker中的scala-library jar包。 - 序列化机制切换:当前使用JavaSerializer,尝试切换为KryoSerializer验证是否为序列化机制导致的类不兼容,在SparkConf中添加
.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")。 - 依赖打包一致性:将本地项目依赖打包为fat jar,通过
spark.jars指定该jar包,确保Driver与Executor使用完全相同的依赖集合,避免两端依赖差异。
内容的提问来源于stack exchange,提问作者Maxim
相关产品推荐
相关产品推荐

