LabHub
学习 学习路径 课程

数据流水线

幂等的 upsert 流水线

在 LabHub 中继续学习

目标

创建一条加载流水线:无论相同输入提交多少次,结果都保持不变,并且只更新内容实际发生变化的行。

为什么重要

流水线必然会失败。网络会中断,进程会在部署期间终止,源系统也可能延迟发送数据。如果每次都需要人工判断“目前已经加载到哪里”,迟早会发生错误。

采用幂等设计后就不再需要这种判断。失败时直接重新运行即可。核心机制有两个:在自然键上设置唯一约束,让数据库阻止重复数据;通过内容哈希判断内容是否真正变化,不触碰值保持不变的行。

第二点尤其重要。如果无条件更新,即使输入没有变化,重复运行也会提升所有行的更新时间。这样就无法判断“究竟发生了哪些实际变化”,下游只获取变化数据的增量消费也会被破坏。

步骤

目标为 staging.orders_raw,规范化规则与前一项实习相同。

  1. 创建 orders_final 表。列依次为 order_ref, order_date, amount, status, content_hash, created_at, updated_at,且 order_ref 必须是主键。时间列的默认值为当前时间。
  2. 创建 v_raw_dedup 视图。每个单据只保留 raw_id 最大的一行;视图应包含 raw_id 列,日期、金额和状态必须是规范化后的值。
  3. v_raw_dedup 加载到 orders_final。如果单据已存在,则更新其值。加载后的记录数必须等于不同单据的数量。
  4. content_hashmd5(coalesce(order_date::text,'') || '|' || coalesce(amount::text,'') || '|' || status) 规则填充。
  5. 再次执行同一项加载,并将运行前后的 orders_final 记录数写入 /root/etl/upsert_rerun.txt每个值单独一行,共两行
  6. staging.orders_raw 中把单据 ORD-000100status 改成另一个值,然后再次执行加载。此时 updated_at 晚于 created_at 的行必须恰好为 1 行
  7. 创建 etl_run_log 表。列依次为 run_id, started_at, inserted_rows, updated_rows;每次运行记录一行,至少应有 3 行。没有变化的运行中,插入数和更新数都必须为 0。
  8. 创建 v_final_check 视图。列为 metric, value,并输出 total_rows, distinct_refs, changed_rows 三行。

参考

使用自然键创建目标表

创建 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_hashmd5(coalesce(order_date::text,'') || '|' || coalesce(amount::text,'') || '|' || status) 规则填充。

按规定顺序拼接各个值后生成哈希。值不存在时使用空字符串。

确认无变化的重复运行

再次执行同一项加载,并将运行前后的 orders_final 记录数写入 /root/etl/upsert_rerun.txt每个值单独一行,共两行

使用相同输入再运行一次,并将运行前后的记录数分两行写入。两个值必须相同。

应用延迟到达的变更

staging.orders_raw 中把单据 ORD-000100status 改成另一个值,然后再次执行加载。此时 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 三行。

以名称和值组成的键值对形式,输出总记录数、不同单据数和已更新行数这三项指标。