You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 13:35:06