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

Java应用通过Spark集群读取Mongo集合时抛出ClassCastException

Spark集群连接MongoDB执行df.show()抛出ClassCastException解决方案

问题现象

宿主机Java应用连接Docker中Spark Master(3.4.1版本),已成功连接MongoDB。执行df.printSchema()可正常输出表结构:

|-- Data_value: double (nullable = true)
 |-- Group: string (nullable = true)
 |-- MAGNTUDE: integer (nullable = true)
 |-- Period: double (nullable = true)
 |-- STATUS: string (nullable = true)
 |-- Series_reference: string (nullable = true)
 |-- Series_title_1: string (nullable = true)
 |-- Subject: string (nullable = true)
 |-- UNITS: string (nullable = true)
 |-- _id: string (nullable = true)

但执行df.show()时抛出如下异常:

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

环境与代码

  • 应用依赖Spark 3.4.1,使用JDK17
  • Spark Master版本3.4.1(运行在Docker容器)
    核心代码:
SparkSession spark = SparkSession.builder()
                .master("spark://localhost:7077")
                .appName("MongoSparkConnectorIntro")
                .config("spark.mongodb.read.connection.uri", "mongodb://127.0.0.1/testdb.balances")
                .config("spark.mongodb.write.connection.uri", "mongodb://127.0.0.1/testdb.balances")
                .getOrCreate();
Dataset<Row> df = spark.read().format("mongodb")
        .load();
df.printSchema();
df.show();

已做排查

  • 将SparkSession的master改为local模式,程序可正常运行
  • 尝试通过spark.config("spark.jars","target/spark-demo.jar")配置添加应用JAR,问题未解决

解决方法

1. 统一依赖版本

该异常本质是Scala集合序列化类型不匹配,核心原因是客户端与Spark集群节点的依赖版本不一致:

  • 确认Docker中Spark集群使用的Scala版本(Spark 3.4.1默认对应Scala 2.12.x),确保应用依赖的Scala版本与其完全一致
  • 选用与Spark 3.4.x兼容的Mongo Spark Connector版本,例如mongo-spark-connector_2.12:10.2.1

2. 正确分发依赖包

Docker内的Spark Worker无法直接访问宿主机本地的JAR文件,需确保依赖包能被所有Worker节点获取:

  • 提交应用时使用--jars参数,将应用JAR、Mongo Connector等依赖包的路径(或远程存储路径)明确指定,让集群自动分发
  • 或者在启动Spark Worker容器时,将本地依赖包目录挂载到容器内的Spark依赖目录

3. 优化序列化配置

添加Kryo序列化配置,避免Scala集合序列化冲突:

SparkSession spark = SparkSession.builder()
        // 保留原有配置
        .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
        .getOrCreate();

4. 修正MongoDB连接地址

Docker容器内的127.0.0.1指向容器自身,需改用宿主机的局域网IP,确保Spark Worker能访问到宿主机上的MongoDB:

.config("spark.mongodb.read.connection.uri", "mongodb://[宿主机局域网IP]:27017/testdb.balances")
.config("spark.mongodb.write.connection.uri", "mongodb://[宿主机局域网IP]:27017/testdb.balances")

内容的提问来源于stack exchange,提问作者help pleh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 07:10:57