LabHub
学习 学习路径 课程

数据流水线

幂等性 — 流水线一定会被再跑一次

在 LabHub 中继续学习

一句话总结

管道失败,失败后重新运行,这时如果结果不同,数据就会被污染,所以重新执行的安全性不是选择,而是设计的前提。

概念图: 让无论放多少次都会得到同样的结果 · 自然键 · 内容哈希 · 需要删除重复

为什么需要这个?

假设部署在凌晨失败了。早上重新转。但是如果失败点是在装载中间的话,一部分已经进去了。就这样重新装载的话会产生重复,如果把整个桶删掉再装载的话,中间进来的其他数据也会被丢失。

如果每次都让别人判断这种情况的话,总有一天会出错。最好一开始就让无论放多少次都会得到同样的结果

怎么行动

特征的起点是自然键。只有每个行唯一标识的值才能识别“已经进入的东西”。比如票据号码、订单号、活动ID等。对这个值施加唯一约束,数据库就会代替防止重复。

接下来就是upsert。

INSERT INTO orders_final (order_ref, amount, status, content_hash)
SELECT ...
ON CONFLICT (order_ref) DO UPDATE
  SET amount = EXCLUDED.amount,
      status = EXCLUDED.status,
      content_hash = EXCLUDED.content_hash,
      updated_at = now()
  WHERE orders_final.content_hash IS DISTINCT FROM EXCLUDED.content_hash;

最后WHERE节很重要。如果没有这个,即使在没有改变内容的重新执行中,所有行的内容也会updated_at这会被更新。那样的话,就无法知道“实际上发生了什么变化”,而且只有下游才会带走变化的部分,导致增量消费也会崩溃。

如果加上内容哈希,比较就会变得简单。将值按确定的顺序连接起来计算哈希,只有当这些值不同时才会更新。即使列变多,比较逻辑也会保持在一行。

需要删除重复。如果同一张票被装两次,需要制定规则,留下哪张。一般留下后来来的。如果有表示装载顺序的增量键在原件上,就可以只选择每张票的最大值行。

按分割槽单元更换

有时不能用行单位upsert。如果原件整块重新给一天的差额, 或者没有办法找出被删除的行的情况。这时,将覆盖单位作为分区 抓住。

BEGIN;
  -- 1) 새 데이터를 임시 표에 적재한다
  CREATE TEMP TABLE stage_20260906 (LIKE orders_final INCLUDING ALL);
  COPY stage_20260906 FROM ...;

  -- 2) 그날 파티션만 통째로 바꿔 끼운다
  ALTER TABLE orders_final DETACH PARTITION orders_20260906;
  ALTER TABLE orders_final ATTACH PARTITION stage_20260906
        FOR VALUES FROM ('2026-09-06') TO ('2026-09-07');
COMMIT;

核心不是**“删除后插入”,而是“制作后留着换上”**。删除后 插入期间会产生没有数据的区间,这时查询的人会看到空的结果。 更换插座是在交易中以原子的方式结束的,所以没有那个部分。

这种方式在重新执行时也很安全。无论转几次,该分区的最终状态都是最后一个 执行结果之一。

识别不安全的运算的方法

不加括号的运算有共同点。读取当前值,以此为基础 是使用的东西

不安全 安全
UPDATE t SET n = n + 1 UPDATE t SET n = <계산된 절대값>
INSERT(无限制) INSERT ... ON CONFLICT DO UPDATE
在文件中append 重新写文件 rename
在Q上发送消息 发送后消费者不会重复发送
外部API调用(支付等) 一起发送账号密码

最后一行很重要。管道呼叫外部系统时,要考虑那边的横向平衡性。 需要期待。大多数结算API是Idempotency-Key收到标题。相同尺寸 叫两次的话,第二次会把第一次的结果原封不动地还给你。即使重新尝试键盘 因为必须相同,所以每个请求不能重新创建,从原始数据中决定性地 必须诱导(例如:sha256(주문번호 + 금액)).

水印和再处理范围

在增量装载中,记录“处理了到哪里”。在处理完这个值后更新 这是原则。首先更新的话,中间死了的话,永远跳过那个部分。

읽기 → 변환 → 적재 → (성공 시에만) 워터마크 갱신

并且在水印上留出余地。以活动时间为基准,原件迟迟未到达的数据 如果允许的话(late arrival),水印不是最后的时间,而是最后的时间-延期期间 抓住。那么每次都会重新读重叠的区间,因为是均匀地设置的。 没有问题。如果有重叠的话,叠加阅读是免费的。

在现场相遇的样子

留下执行日志会很有帮助。每次执行时记录插入次数和更新次数,没有变化的执行都会保持0,晚到的更改反映在执行中的更新次数会增加。只要看这两个数字,就可以判断管道是否正常,源头是否存在异常。

而且重试代码中只需要放入数据库操作。如果混合了发送邮件或外部API调用,每次重试都会重复其副作用。如果需要外部调用,那边也应该是收到偏移量键的方式。

下次实习要做的事情

在自然键上挂上基本键,删除重复表格,并通过内容哈希判断是否更改的upsert。然后确认没有变化的重新执行和只更改了一件的重新执行分别如何记录。