如何将MongoDB文档以单文本列读入PySpark DataFrame(无模式)
无模式读取MongoDB文档到Databricks Delta Lake青铜层
核心思路
绕过PySpark MongoDB连接器的模式推断限制,将整个MongoDB文档打包为单一字符串字段(value),同时添加时间戳字段(time_stamp),完全贴合青铜层的无模式存储需求。
解决方案1:转为JSON字符串(推荐,兼容性好)
这种方式将原始文档序列化为标准JSON字符串,无需额外依赖,后续处理更灵活。
# 1. 读取MongoDB集合,关闭模式推断 df_raw = spark.read.format("mongodb") \ .option("spark.mongodb.input.uri", "mongodb://<host>:<port>/<your_db>.<your_collection>") \ .option("spark.mongodb.input.schemaInference", "false") \ .load() # 2. 将整个文档转为JSON字符串作为value列,添加当前时间戳 from pyspark.sql.functions import to_json, struct, current_timestamp df_bronze = df_raw.select( to_json(struct("*")).alias("value"), # struct("*")打包所有字段为单一结构体,再转JSON current_timestamp().alias("time_stamp") ) # 3. 写入青铜层Delta表 df_bronze.write.format("delta") \ .mode("append") \ .save("/dbfs/path/to/bronze/your_table")
解决方案2:保留原始BSON格式(二进制转Base64)
如果需要严格保留BSON的原始格式(比如包含JSON不支持的类型),可以将BSON二进制数据编码为Base64字符串存储。
# 1. 定义UDF:将文档转为BSON并编码为Base64字符串 import base64 import bson from pyspark.sql.functions import udf from pyspark.sql.types import StringType def bson_to_base64(doc): bson_bytes = bson.BSON.encode(doc.asDict()) # 将Row转为字典后编码为BSON return base64.b64encode(bson_bytes).decode("utf-8") bson_base64_udf = udf(bson_to_base64, StringType()) # 2. 读取MongoDB并处理 df_raw = spark.read.format("mongodb") \ .option("spark.mongodb.input.uri", "mongodb://<host>:<port>/<your_db>.<your_collection>") \ .option("spark.mongodb.input.schemaInference", "false") \ .load() df_bronze = df_raw.select( bson_base64_udf(struct("*")).alias("value"), current_timestamp().alias("time_stamp") ) # 3. 写入Delta表 df_bronze.write.format("delta") \ .mode("append") \ .save("/dbfs/path/to/bronze/your_table")
可选优化:使用MongoDB ObjectId生成时间戳
如果希望时间戳与文档创建时间一致(而非读取时间),可以从MongoDB的_id(ObjectId)中提取创建时间:
from pyspark.sql.functions import expr df_bronze = df_raw.select( to_json(struct("*")).alias("value"), # 从ObjectId的前8位十六进制字符串解析出创建时间 expr("toTimestamp(from_unixtime(conv(substring(_id, 1, 8), 16, 10)))").alias("time_stamp") )
关键注意事项
- 必须设置
spark.mongodb.input.schemaInference=false,禁止Spark自动推断模式,确保原始文档结构完全保留 struct("*")会将DataFrame的所有字段打包为一个结构体,保证整个文档被作为单一值处理- 写入Delta表时建议使用
append模式,符合青铜层增量写入的特性
内容的提问来源于stack exchange,提问作者Shay Palachy Affek
相关产品推荐
相关产品推荐

