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
相关产品推荐
相关产品推荐

