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)
问题原因
- 增量数据写入错误:创建增量表时,代码误使用
pd.DataFrame(data)而非pd.DataFrame(data_updated),导致增量表中存储的是全量数据,而非预期的更新行。 - 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
相关产品推荐
相关产品推荐

