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

AWS Glue中Delta表MERGE INTO未更新反而插入行问题排查

Delta Lake MERGE INTO 未更新反而新增行的问题排查

问题背景

我正尝试搭建一个演示数据湖框架Upsert功能的AWS Glue作业:

  • 全量数据已保存为S3存储桶中的Delta表:
data = {'visitor': ['foo', 'bar', 'baz'], 
         'id': [1, 2, 3],
         'B': [1, 0, 1],
         'C': [1, 0, 0]}  
  • 增量数据计划存储在S3不同前缀下:
data_updated = {'visitor': ['foo_updated'], 
        'id': [1],
        'B': [1],
        'C': [1]} 

执行MERGE INTO语句后,id=1的增量行未更新全量表原有行,反而被追加插入。

执行的MERGE语句

delta_df = DeltaTable.forPath(spark, "s3://example_bucket/full_load")
cdc_df = spark.read.format("delta").load("s3://example_bucket/incremental_load/")

final_df = delta_df.alias("prev_df").merge( \
source = cdc_df.alias("append_df"), \
#matching on primarykey
condition = expr("prev_df.id = append_df.id"))\
.whenMatchedUpdate(set= {
    "prev_df.B"           : col("append_df.B"), 
    "prev_df.C"         : col("append_df.C"),
    "prev_df.visitor"        : col("append_df.visitor")} )\
.whenNotMatchedInsert(values =
#inserting a new row to Delta table
{   "prev_df.B"             : col("append_df.B"),
    "prev_df.C"           : col("append_df.C"), 
    "prev_df.visitor"         : col("append_df.visitor"),
})\
.execute()

表的创建代码

df = pd.DataFrame(data)
dataFrame = spark.createDataFrame(df)
dataFrame.write \
    .format("delta") \
    .mode("overwrite") \
    .save("s3://example_bucket/full_load", overwrite=True)

df = pd.DataFrame(data)
dataFrame = spark.createDataFrame(df)
dataFrame.write \
    .format("delta") \
    .mode("overwrite") \
    .save("s3://example_bucket/incremental_load", overwrite=True)

问题原因

  1. 增量数据写入错误:创建增量表时,代码误使用pd.DataFrame(data)而非pd.DataFrame(data_updated),导致增量表中存储的是全量数据,而非预期的更新行。
  2. MERGE语句字段名格式错误:whenMatchedUpdate和whenNotMatchedInsert中的字段键不应该带表别名前缀(prev_df.),Delta Lake要求直接使用目标表的字段名,无需添加别名。

修复方案

1. 修正增量表写入代码

替换增量数据的数据源:

df = pd.DataFrame(data_updated)  # 使用正确的增量数据集
dataFrame = spark.createDataFrame(df)
dataFrame.write \
    .format("delta") \
    .mode("overwrite") \
    .save("s3://example_bucket/incremental_load", overwrite=True)

2. 修正MERGE语句的字段映射

去掉字段名中的prev_df.前缀:

delta_df = DeltaTable.forPath(spark, "s3://example_bucket/full_load")
cdc_df = spark.read.format("delta").load("s3://example_bucket/incremental_load/")

delta_df.alias("prev_df").merge( \
source = cdc_df.alias("append_df"), \
condition = expr("prev_df.id = append_df.id"))\
.whenMatchedUpdate(set= {
    "B"           : col("append_df.B"), 
    "C"         : col("append_df.C"),
    "visitor"        : col("append_df.visitor")} )\
.whenNotMatchedInsert(values =
{   "B"             : col("append_df.B"),
    "C"           : col("append_df.C"), 
    "visitor"         : col("append_df.visitor")
})\
.execute()

验证逻辑

修复后,MERGE INTO会根据id匹配全量表中的行:

  • 匹配到id=1的行时,会用增量数据更新visitor、B、C字段
  • 若增量数据中有新id,会执行插入操作

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 15:55:25