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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 11:27:20