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

在AWS Glue中通过PySpark加载MongoDB Atlas数据失败求助

AWS Glue加载MongoDB Atlas类型不匹配问题解决方案

针对你遇到的metadata.dateButtonInstalled字段字符串类型被错误推断为timestamp导致的报错,以下是无需修改源数据的可行配置方案:

方案1:自定义Schema强制指定字段类型

关闭自动类型推断,手动定义Spark Schema明确字段类型,避免连接器误判:

from pyspark.sql.types import StructType, StructField, StringType

# 根据你的集合实际结构编写Schema,此处仅为示例
custom_schema = StructType([
    StructField("metadata", StructType([
        StructField("dateButtonInstalled", StringType(), nullable=True),
        # 补充其他字段的类型定义
    ])),
    # 添加顶层字段的类型定义
])

df = (
    self.spark.read.format("mongodb")
    .option("connection.uri", "mongodb+srv://<something>.<some-thing>.mongodb.net/<db>?authSource=<user>")
    .option("collection", "你的目标集合名")  # 填写实际集合名称
    .option("spark.mongodb.read.schema.infer.enabled", "false")
    .schema(custom_schema)
    .load()
    .limit(5)
)

后续可通过Spark函数将字符串日期转换为timestamp类型:

from pyspark.sql.functions import to_timestamp

df = df.withColumn(
    "metadata.dateButtonInstalled",
    to_timestamp("metadata.dateButtonInstalled", "yyyy-MM-dd'T'HH:mm:ss.SSS'Z'")
)

方案2:禁用日期自动解析为Timestamp

配置连接器不将字符串格式的日期推断为timestamp,强制以字符串类型读取:

df = (
    self.spark.read.format("mongodb")
    .option("connection.uri", "mongodb+srv://<something>.<some-thing>.mongodb.net/<db>?authSource=<user>")
    .option("collection", "你的目标集合名")
    .option("spark.mongodb.read.schema.preferTimestamp", "false")
    .option("spark.mongodb.read.dateFormat", "none")
    .load()
    .limit(5)
)

方案3:增大类型推断采样量

若自动推断时采样的文档不具备代表性(如采样文档中该字段为timestamp,但实际存在字符串值),可增大采样样本量提升推断准确性(注意:过大样本量会影响读取性能):

df = (
    self.spark.read.format("mongodb")
    .option("connection.uri", "mongodb+srv://<something>.<some-thing>.mongodb.net/<db>?authSource=<user>")
    .option("collection", "你的目标集合名")
    .option("spark.mongodb.read.schema.sampleSize", "5000")  # 调整为合适的样本数量,默认1000
    .load()
    .limit(5)
)

方案4:忽略类型不匹配的畸形记录

若允许跳过少量类型异常的记录,可启用忽略畸形记录配置:

df = (
    self.spark.read.format("mongodb")
    .option("connection.uri", "mongodb+srv://<something>.<some-thing>.mongodb.net/<db>?authSource=<user>")
    .option("collection", "你的目标集合名")
    .option("spark.mongodb.read.ignoreMalformedRecords", "true")
    .load()
    .limit(5)
)

内容的提问来源于stack exchange,提问作者RushHour

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 16:47:04