批量装载与水位线
目标
亲手运行批处理流水线的一个完整周期,包括提取、加载、记录水位线、增量处理和验证重复运行结果。
为什么重要
批处理流水线发生事故的位置几乎是固定的:边界条件和重复运行。
边界条件中,> 与 >= 仅一个符号的差异,就会导致每次运行都重复或遗漏一条数据。如果批处理每天运行一次,一年就会有 365 条数据出错,而且期间可能无人察觉。因此,最好在流水线中加入提取后立即将记录数与源数据核对的验证步骤。
重复运行更为重要。流水线必然会失败,失败后就需要重新运行。如果重新运行后的结果发生变化,数据就会被污染。因此,必须采用这样的设计:在加载目标表上设置唯一约束,并在发生冲突时忽略或更新记录。本实习将同一项增量作业运行两次,亲自确认记录数是否保持不变。
步骤
工作目录为 /root/etl。
- 提取
ordered_at位于 2025 年 1 月的订单,将其id,ordered_at,status,total_amount写入/root/etl/orders_2025_01.csv。包含表头。时间范围为大于或等于TIMESTAMPTZ '2025-01-01 00:00:00+09',且小于TIMESTAMPTZ '2025-02-01 00:00:00+09'。 - 确认文件行数等于该时间范围内的订单数加上一行表头。
- 创建
orders_archive表。列依次为id,ordered_at,status,total_amount,并且id必须具有主键或唯一约束。 - 将第 1 步生成的 CSV 加载到
orders_archive。 - 创建
etl_watermark表。列为job_name,last_ordered_at;在job_name为orders_archive的行中记录已加载数据的最大ordered_at。 - 增量加载水位线之后、
TIMESTAMPTZ '2025-03-01 00:00:00+09'之前的订单,并更新水位线。加载完成后,orders_archive应同时包含 1 月和 2 月的订单。 - 再次执行相同的增量作业。将运行前后的
orders_archive记录数写入/root/etl/rerun.txt,每个值单独一行,共两行。两个值必须相同,并且不得存在重复行。 - 将
orders_archive按状态统计的记录数保存到/root/etl/report.tsv。状态与记录数之间使用制表符分隔,且不写入表头。
参考
\copy (SELECT ...) TO '/root/etl/x.csv' WITH (FORMAT csv, HEADER true)\copy 표이름 (컬럼목록) FROM '/root/etl/x.csv' WITH (FORMAT csv, HEADER true)- 忽略冲突:
INSERT ... ON CONFLICT (id) DO NOTHING - 常见错误 1:两侧边界值都包含时,月末订单会被加载两次。
- 常见错误 2:增量加载后遗漏水位线更新,导致下次运行再次读取同一区间。
将 1 月订单提取为 CSV
提取 ordered_at 位于 2025 年 1 月的订单,将其 id, ordered_at, status, total_amount 写入 /root/etl/orders_2025_01.csv。包含表头。时间范围为大于或等于 TIMESTAMPTZ '2025-01-01 00:00:00+09',且小于 TIMESTAMPTZ '2025-02-01 00:00:00+09'。
psql 的 \copy 会写入客户端一侧的文件。它提供包含表头的选项。
验证提取记录数
确认文件行数等于该时间范围内的订单数加上一行表头。
文件行数应等于数据行数加上一行表头。请确认边界条件只包含一侧。
创建加载目标表
创建 orders_archive 表。列依次为 id, ordered_at, status, total_amount,并且 id 必须具有主键或唯一约束。
重复运行的安全性在这里决定。需要设置约束,确保同一订单不能写入两次。
将 CSV 加载到表中
将第 1 步生成的 CSV 加载到 orders_archive。
\copy 也可以反向使用。不要忘记跳过表头行的选项。
记录水位线
创建 etl_watermark 表。列为 job_name, last_ordered_at;在 job_name 为 orders_archive 的行中记录已加载数据的最大 ordered_at。
在表中记录处理进度。该值必须是已加载数据的最后时间。
增量加载 2 月数据
增量加载水位线之后、TIMESTAMPTZ '2025-03-01 00:00:00+09' 之前的订单,并更新水位线。加载完成后,orders_archive 应同时包含 1 月和 2 月的订单。
只筛选并加载水位线之后的数据,完成后再更新水位线。
确认运行两次仍得到相同结果
再次执行相同的增量作业。将运行前后的 orders_archive 记录数写入 /root/etl/rerun.txt,每个值单独一行,共两行。两个值必须相同,并且不得存在重复行。
再次执行同一项增量作业,并将运行前后的记录数分两行写入文件。需要使用发生冲突时忽略记录的子句。
生成加载结果报告
将 orders_archive 按状态统计的记录数保存到 /root/etl/report.tsv。状态与记录数之间使用制表符分隔,且不写入表头。
将各状态的记录数用制表符分隔后保存。只保留数值,不要写入表头。