幂等的 upsert 流水线
目标
创建一条加载流水线:无论相同输入提交多少次,结果都保持不变,并且只更新内容实际发生变化的行。
为什么重要
流水线必然会失败。网络会中断,进程会在部署期间终止,源系统也可能延迟发送数据。如果每次都需要人工判断“目前已经加载到哪里”,迟早会发生错误。
采用幂等设计后就不再需要这种判断。失败时直接重新运行即可。核心机制有两个:在自然键上设置唯一约束,让数据库阻止重复数据;通过内容哈希判断内容是否真正变化,不触碰值保持不变的行。
第二点尤其重要。如果无条件更新,即使输入没有变化,重复运行也会提升所有行的更新时间。这样就无法判断“究竟发生了哪些实际变化”,下游只获取变化数据的增量消费也会被破坏。
步骤
目标为 staging.orders_raw,规范化规则与前一项实习相同。
- 创建
orders_final表。列依次为order_ref,order_date,amount,status,content_hash,created_at,updated_at,且order_ref必须是主键。时间列的默认值为当前时间。 - 创建
v_raw_dedup视图。每个单据只保留raw_id最大的一行;视图应包含raw_id列,日期、金额和状态必须是规范化后的值。 - 将
v_raw_dedup加载到orders_final。如果单据已存在,则更新其值。加载后的记录数必须等于不同单据的数量。 content_hash按md5(coalesce(order_date::text,'') || '|' || coalesce(amount::text,'') || '|' || status)规则填充。- 再次执行同一项加载,并将运行前后的
orders_final记录数写入/root/etl/upsert_rerun.txt,每个值单独一行,共两行。 - 在
staging.orders_raw中把单据ORD-000100的status改成另一个值,然后再次执行加载。此时updated_at晚于created_at的行必须恰好为 1 行。 - 创建
etl_run_log表。列依次为run_id,started_at,inserted_rows,updated_rows;每次运行记录一行,至少应有 3 行。没有变化的运行中,插入数和更新数都必须为 0。 - 创建
v_final_check视图。列为metric,value,并输出total_rows,distinct_refs,changed_rows三行。
参考
- 冲突时更新:
INSERT ... ON CONFLICT (order_ref) DO UPDATE SET ... WHERE 대상.content_hash IS DISTINCT FROM EXCLUDED.content_hash - 区分插入与更新:可在 CTE 中接收并统计
RETURNING (xmax = 0) AS inserted。 - 每个单据的最新行:
SELECT DISTINCT ON (order_ref) * FROM ... ORDER BY order_ref, raw_id DESC - 常见错误 1:没有为
DO UPDATE添加条件,导致无变化的重复运行也会提升所有行的更新时间。 - 常见错误 2:未经去重就加载,导致同一单据在一个批次中出现两次并引发错误。
使用自然键创建目标表
创建 orders_final 表。列依次为 order_ref, order_date, amount, status, content_hash, created_at, updated_at,且 order_ref 必须是主键。时间列的默认值为当前时间。
将单据编号设为主键,同时设置创建时间和更新时间列。
删除重复单据
创建 v_raw_dedup 视图。每个单据只保留 raw_id 最大的一行;视图应包含 raw_id 列,日期、金额和状态必须是规范化后的值。
同一单据到达两次时保留后到的数据。PostgreSQL 的 DISTINCT ON 适合完成此任务。
创建冲突时更新的加载操作
将 v_raw_dedup 加载到 orders_final。如果单据已存在,则更新其值。加载后的记录数必须等于不同单据的数量。
为 INSERT 添加冲突处理子句。加载前必须规范化日期和状态。
计算内容哈希
content_hash 按 md5(coalesce(order_date::text,'') || '|' || coalesce(amount::text,'') || '|' || status) 规则填充。
按规定顺序拼接各个值后生成哈希。值不存在时使用空字符串。
确认无变化的重复运行
再次执行同一项加载,并将运行前后的 orders_final 记录数写入 /root/etl/upsert_rerun.txt,每个值单独一行,共两行。
使用相同输入再运行一次,并将运行前后的记录数分两行写入。两个值必须相同。
应用延迟到达的变更
在 staging.orders_raw 中把单据 ORD-000100 的 status 改成另一个值,然后再次执行加载。此时 updated_at 晚于 created_at 的行必须恰好为 1 行。
修改源数据中的一份单据后再次运行。必须添加条件,避免连值未变化的行也被更新。
记录运行日志
创建 etl_run_log 表。列依次为 run_id, started_at, inserted_rows, updated_rows;每次运行记录一行,至少应有 3 行。没有变化的运行中,插入数和更新数都必须为 0。
每次运行都记录插入数和更新数。没有任何变化的运行中,两者都为 0。
创建检查指标视图
创建 v_final_check 视图。列为 metric, value,并输出 total_rows, distinct_refs, changed_rows 三行。
以名称和值组成的键值对形式,输出总记录数、不同单据数和已更新行数这三项指标。