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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 17:56:21