在Databricks用PySpark实现Delta Lake去重,仅保留最新crt_dt记录
解决方案:Delta Lake去重并保留最新记录(按product_code+date维度)
针对你描述的每月重复加载未来6个月预测数据、需要按product_code和date保留crt_dt最新记录的场景,结合Delta Lake的特性,以下是几个可行的解决方案:
1. 增量加载首选:MERGE INTO 事务性去重
这是增量更新场景下最高效的方式,每次加载新数据时,仅对需要更新或插入的记录操作,同时利用crt_dt分区减少扫描量。
MERGE INTO target_delta_table t USING new_prediction_data s ON t.product_code = s.product_code AND t.date = s.date WHEN MATCHED AND t.crt_dt < s.crt_dt THEN UPDATE SET quantity = s.quantity, crt_dt = s.crt_dt WHEN NOT MATCHED THEN INSERT (product_code, date, quantity, crt_dt) VALUES (s.product_code, s.date, s.quantity, s.crt_dt)
优势:
- 原子性操作,保证数据一致性
- 自动利用分区 pruning,只扫描与新数据匹配的分区,性能高效
- 仅更新旧于新数据的记录,避免不必要的写入
2. 批量历史数据清洗:窗口函数筛选后覆盖写入
如果需要一次性清理历史重复数据,或者定期执行全量去重,可以用窗口函数筛选出每个product_code+date组内的最新记录,再覆盖写入原表。
WITH latest_records AS ( SELECT product_code, date, quantity, crt_dt, -- 按分组降序排序,取第一条即为最新记录 ROW_NUMBER() OVER (PARTITION BY product_code, date ORDER BY crt_dt DESC) AS rn FROM target_delta_table -- 可选:如果数据量极大,可添加分区过滤,比如只处理近12个月的分区 -- WHERE crt_dt >= DATE_SUB(CURRENT_DATE(), 365) ) INSERT OVERWRITE target_delta_table SELECT product_code, date, quantity, crt_dt FROM latest_records WHERE rn = 1
注意:覆盖写入会替换整个表(或指定分区),建议先在测试环境验证,或结合分区过滤缩小操作范围。
3. 简化场景:Delete+Insert 组合操作
如果能保证每月加载的新数据crt_dt一定是当前最新值(比如crt_dt设为加载当月的第一天/最后一天,且每月只加载一次),可以先删除对应product_code+date的旧记录,再插入新数据,Delta Lake会保证这个操作的原子性。
-- 删除与新数据重复的旧记录 DELETE FROM target_delta_table WHERE (product_code, date) IN (SELECT product_code, date FROM new_prediction_data) -- 插入最新的预测数据 INSERT INTO target_delta_table SELECT product_code, date, quantity, crt_dt FROM new_prediction_data
优势:逻辑简单,适合确定性的最新数据场景,结合分区能快速定位待删除的记录。
额外优化建议
- 确保
crt_dt为日期/时间类型而非字符串,避免排序和比较出错 - 定期执行
OPTIMIZE target_delta_table ZORDER BY (product_code, date),优化按这两个维度查询和更新的性能 - 增量加载时,可对新数据先进行过滤(比如只保留未来6个月的
date),减少不必要的写入
内容的提问来源于stack exchange,提问作者Krishna
相关产品推荐
相关产品推荐

