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

AWS Glue转换Timestamp写入PostgreSQL触发非空约束错误排查

AWS Glue 转换CSV时间字段到Timestamp时触发PostgreSQL非空约束错误

我在做AWS Glue任务POC,从S3读取CSV文件转换Schema后写入PostgreSQL。当所有字段设为string时数据正常读取,但把load_dt设为timestamp类型时,出现以下错误:

"An error occurred while calling o 120.pyWriteDynamicFrame. ERROR: null value column load_dt of relation t002_stg violates not-null constraint"

已经确认CSV里仅有的两条数据都包含该时间戳字段,不想用转成PySpark DataFrame再处理的繁琐方式。相关代码如下:

AmazonS3_node12345 = glueContext.create_dynamic_frame.from_options(
    format_options={"quoteChar": '"', "withHeader": True, "separator": "|"},
    connection_type="s3",
    format="csv",
    connection_options={
        "paths": ["s3://dedicated_drop_folder/decrypted/product_lookups.csv"],
        "recurse": True,
    },
    transformation_ctx="AmazonS3_node12345",
)

# Script generated for node Change Schema
ChangeSchema_node12345 = ApplyMapping.apply(
    frame=AmazonS3_node12345,
    mappings=[
        ("product_name", "string", "product_name", "string"),
        ("pd_parent_name", "string", "pd_parent_name", "string"),
        ("LAST_LOAD_DT", "string", "LOAD_DT", "timestamp"),
        ("DATA_SRC_NM", "string", "DATA_SRC_NM", "string"),
    ],
    transformation_ctx="ChangeSchema_node12345",
)

# Script generated for node PostgreSQL
PostgreSQL_node12345 = glueContext.write_dynamic_frame.from_catalog(
    frame=ChangeSchema_node12345,
    database="dev_testing",
    table_name="crm_pricing_lk_stg",
    transformation_ctx="PostgreSQL_node12345",
)

job.commit()

问题根源

ApplyMapping默认的timestamp转换逻辑无法匹配CSV中LAST_LOAD_DT的时间格式,导致转换失败,字段值变成null,触发了PostgreSQL的非空约束。

无需转DataFrame的解决方案

1. 读取CSV时直接指定时间格式

在读取CSV的format_options里添加timestampFormat,明确时间字符串的格式,让Glue在读取阶段就把LAST_LOAD_DT解析成timestamp类型,后续ApplyMapping直接映射同类型即可:

AmazonS3_node12345 = glueContext.create_dynamic_frame.from_options(
    format_options={
        "quoteChar": '"', 
        "withHeader": True, 
        "separator": "|",
        "timestampFormat": "yyyy-MM-dd HH:mm:ss.SSS"  # 替换成你CSV里实际的时间格式
    },
    connection_type="s3",
    format="csv",
    connection_options={
        "paths": ["s3://dedicated_drop_folder/decrypted/product_lookups.csv"],
        "recurse": True,
    },
    transformation_ctx="AmazonS3_node12345",
)

# ApplyMapping中直接映射timestamp类型
ChangeSchema_node12345 = ApplyMapping.apply(
    frame=AmazonS3_node12345,
    mappings=[
        ("product_name", "string", "product_name", "string"),
        ("pd_parent_name", "string", "pd_parent_name", "string"),
        ("LAST_LOAD_DT", "timestamp", "LOAD_DT", "timestamp"),
        ("DATA_SRC_NM", "string", "DATA_SRC_NM", "string"),
    ],
    transformation_ctx="ChangeSchema_node12345",
)

2. 用ResolveChoice处理转换失败的字段

如果时间格式不固定,或者读取时没指定格式导致转换失败,可以用ResolveChoice强制转换字段类型,避免生成null:

from awsglue.transforms import ResolveChoice

# 在ApplyMapping之后添加ResolveChoice节点,指定LOAD_DT字段强制转为timestamp
ResolveChoice_node = ResolveChoice.apply(
    frame=ChangeSchema_node12345,
    specs=[("LOAD_DT", "cast:timestamp")],
    transformation_ctx="ResolveChoice_node",
)

# 再写入PostgreSQL
PostgreSQL_node12345 = glueContext.write_dynamic_frame.from_catalog(
    frame=ResolveChoice_node,
    database="dev_testing",
    table_name="crm_pricing_lk_stg",
    transformation_ctx="PostgreSQL_node12345",
)

3. 用Map转换手动处理时间解析

通过Map转换自定义处理逻辑,直接在DynamicFrame层面解析时间字段:

from awsglue.transforms import Map
from pyspark.sql.functions import to_timestamp

def parse_time(record):
    # 按照CSV实际格式解析时间,替换成你的格式字符串
    record["LOAD_DT"] = to_timestamp(record["LAST_LOAD_DT"], "yyyy-MM-dd HH:mm:ss.SSS")
    return record

# 先处理时间字段,再做字段映射
Mapped_node = Map.apply(
    frame=AmazonS3_node12345,
    f=parse_time,
    transformation_ctx="Mapped_node",
)

ChangeSchema_node12345 = ApplyMapping.apply(
    frame=Mapped_node,
    mappings=[
        ("product_name", "string", "product_name", "string"),
        ("pd_parent_name", "string", "pd_parent_name", "string"),
        ("LOAD_DT", "timestamp", "LOAD_DT", "timestamp"),
        ("DATA_SRC_NM", "string", "DATA_SRC_NM", "string"),
    ],
    transformation_ctx="ChangeSchema_node12345",
)

验证方法

转换后可以通过以下方式确认LOAD_DT字段是否正常:

# 查看Schema
ChangeSchema_node12345.printSchema()
# 查看样本数据
ChangeSchema_node12345.toDF().limit(10).show()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 11:22:20