PySpark读写MongoDB时如何保留_id的ObjectId类型
解决PySpark写入MongoDB时保持_id为ObjectId类型的问题
可行方案步骤
转换_id为标准结构体
PySpark读取MongoDB后,_id会被转为字符串格式,需要将其重构为MongoDB连接器能识别的ObjectId结构体(包含oid字段):from pyspark.sql.functions import struct, col # 把字符串类型的_id转成oid结构体 df = df.withColumn("_id", struct(col("_id").alias("oid")))配置Spark连接器参数
创建SparkSession时,添加必要的配置,确保连接器能正确解析结构体为ObjectId:from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("MongoDBWriteWithObjectId") \ .config("spark.mongodb.input.uri", "mongodb://你的主机:27017/你的数据库.你的集合") \ .config("spark.mongodb.output.uri", "mongodb://你的主机:27017/你的数据库.你的集合") \ .config("spark.mongodb.output.replaceDocument", "true") # 确保更新现有文档而非新增 .getOrCreate()选择正确的写入模式
避免用append模式生成重复的字符串_id文档,更新用update模式,覆盖用overwrite模式:df.write \ .format("mongodb") \ .mode("update") \ .save()
之前方法无效的原因
- 直接用UDF返回
bson.ObjectId对象:PySpark无法序列化bson原生对象到MongoDB,必须转为连接器兼容的结构体格式。 - 设置
convertJson为objectOrArrayOnly:该选项仅控制JSON转换范围,对ObjectId的识别无帮助,保持默认none即可。
内容的提问来源于stack exchange,提问作者Arno
相关产品推荐
相关产品推荐

