Spark DataFrame写入Mongo失败求助:空指针异常问题
解决MongoDB Spark Connector 10.0.1写入时的NullPointerException问题
问题场景
- 使用MongoDB Spark Connector版本:10.0.1
- 写入配置:
.config("spark.mongodb.write.connection.uri", "mongodb://127.0.0.1:27017/") .config("spark.mongodb.write.database", "test") .config("spark.mongodb.write.collection", "spark") .config("spark.mongodb.write.operationType", "insert")
- 执行写入代码:
spark_df.write.format("mongodb").mode("append").save()
- 抛出异常:
24/01/09 11:17:26 ERROR Utils: Aborting task java.lang.NullPointerException at org.apache.spark.api.python.SerDeUtil$.$anonfun$toJavaArray$1(SerDeUtil.scala:72) ...(省略重复日志内容)
已尝试更换连接器版本、匹配Spark官方文档要求的版本,问题仍未解决。
排查方向与解决方案
1. 检查DataFrame中的空值与特殊数组类型
异常触发在Spark的序列化工具类中,大概率是DataFrame存在无法被Mongo连接器序列化的空值或特殊数组:
- 比如Python DataFrame中包含
None的数组列(ArrayType(NullType)类型) - 混合数据类型的数组(同时包含数值、字符串的数组)
先打印DataFrame的Schema确认字段类型:
spark_df.printSchema()
如果发现可疑字段,尝试清理或转换:
// 过滤包含空数组的行 val cleaned_df = spark_df.filter(!col("target_array_col").isNull) // 或把空数组替换为默认值 val cleaned_df = spark_df.withColumn("target_array_col", when(col("target_array_col").isNull, array(lit(""))).otherwise(col("target_array_col")))
2. 清理冲突依赖
Spark环境中可能存在旧版Mongo Java Driver或序列化相关的冲突依赖:
- 清理Spark安装目录下
jars文件夹中的无关Mongo依赖 - 使用
--packages参数提交任务,让Spark自动管理依赖版本:
spark-submit --packages org.mongodb.spark:mongo-spark-connector_2.12:10.0.1 your_script.py
3. 调整写入参数逻辑
当前配置同时指定了operationType=insert和mode=append,二者存在逻辑重叠,可能触发异常:
- 移除
spark.mongodb.write.operationType配置,仅保留mode=append:
.config("spark.mongodb.write.connection.uri", "mongodb://127.0.0.1:27017/") .config("spark.mongodb.write.database", "test") .config("spark.mongodb.write.collection", "spark")
4. 测试最小可复现案例
用极简DataFrame测试写入,排除业务数据的影响:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.types._ val spark = SparkSession.builder() .appName("MongoTest") .config("spark.mongodb.write.connection.uri", "mongodb://127.0.0.1:27017/") .config("spark.mongodb.write.database", "test") .config("spark.mongodb.write.collection", "spark") .getOrCreate() val testData = Seq( (1, "test1", Array(1,2,3)), (2, "test2", Array(4,5,6)) ) val test_df = spark.createDataFrame(testData).toDF("id", "name", "nums") test_df.write.format("mongodb").mode("append").save()
如果该案例成功,说明问题出在业务DataFrame的数据类型上,需针对性排查。
内容的提问来源于stack exchange,提问作者sarang pratham
相关产品推荐
相关产品推荐

