冪等性 — パイプラインは必ずまた回ることになる
一言でいうと
パイプラインは失敗し、失敗したらもう一度実行し、そのとき結果が変わるとデータが汚染されるため、再実行の安全性は、選択ではなく設計の前提です。
なぜ必要なのか
バッチが明け方に失敗したとします。朝にもう一度実行します。ところが、失敗した地点がロードの途中だったなら、一部はすでに入っています。そのままもう一度入れると重複が生じ、丸ごと消してからもう一度入れると、その間に入ってきた別のデータまで消えます。
この状況を、毎回人に判断させると、いつか間違いが起きます。そもそも何回入れても同じ結果になるように作るほうがよいです。
どう動くのか
冪等性の出発点は、自然キーです。各行を一意に識別する値があって初めて、「すでに入ってきたもの」を見分けられます。伝票番号、注文番号、イベント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句が重要です。これがないと、内容が1つも変わっていない再実行でも、すべての行のupdated_atが更新されます。そうなると、「何が実際に変わったか」がわからなくなり、下流で変更分だけを取得する増分の消費も壊れます。
内容ハッシュを置くと、比較が簡単になります。値を決まった順序でつなげてハッシュを計算し、その値が異なるときだけ更新します。カラムが増えても、比較のロジックは1行のままです。
重複排除も必要です。同じ伝票が2回ロードされたなら、どちらを残すか、ルールを決める必要があります。普通は、あとから入ってきたものを残します。元データにロード順序を表す増加キーがあれば、伝票ごとに最大値の行だけを選べばよいです。
パーティション単位で入れ替える
行単位のupsertではうまくいかない場合があります。元データが1日分を丸ごとまた渡してくる場合や、削除された行を見つける方法がない場合です。このときは、上書きの単位をパーティションにします。
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;
核心は、消して入れるのではなく、作っておいて差し替えることです。消して入れる間には、データがない区間ができ、そのときに参照した人は、空の結果を見ます。差し替えは、トランザクションの中でアトミックに終わるため、その区間がありません。
この方式は、再実行にも安全です。何回回しても、そのパーティションの最終状態は、最後の実行結果1つです。
安全でない操作の見分け方
冪等でない操作には、共通点があります。現在の値を読んで、それをもとに書き込むことです。
| 安全でない | 安全 |
|---|---|
UPDATE t SET n = n + 1 |
UPDATE t SET n = <계산된 절대값>(プレースホルダーは計算済みの絶対値です) |
INSERT(制約なし) |
INSERT ... ON CONFLICT DO UPDATE |
| ファイルにappend | ファイルを新しく書いてrename |
| キューにメッセージを発行 | 発行するが、コンシューマーが重複を除く |
| 外部API呼び出し(決済など) | 冪等キーを一緒に送る |
最後の行が重要です。パイプラインが外部システムを呼ぶときは、その側の冪等性に頼る必要があります。ほとんどの決済APIは、Idempotency-Keyヘッダーを受け取ります。同じキーで2回呼ぶと、2回目は最初の結果をそのまま返します。キーは、リトライしても同じである必要があるため、リクエストごとに新しく作ってはいけず、元データから決定論的に導く必要があります(例はsha256(주문번호 + 금액)で、プレースホルダーは注文番号と金額です)。
ウォーターマークと再処理の範囲
増分ロードでは、「どこまで処理したか」を記録します。この値を処理したあとに更新するのが原則です。先に更新すると、途中で落ちたときに、その区間を永遠に飛ばします。
읽기 → 변환 → 적재 → (성공 시에만) 워터마크 갱신
そして、ウォーターマークに余裕を持たせます。元データが、イベント時刻基準で遅れて届くデータを許容するなら(late arrival)、ウォーターマークを、最後の時刻ではなく最後の時刻−猶予期間にします。そうすると、重なる区間を毎回読み直すことになりますが、冪等に作ってあるので、問題になりません。冪等性があれば、重ねて読むのはただ同然になります。
現場での姿
実行ログを残すと、大きな助けになります。実行ごとに挿入件数と更新件数を記録しておけば、何も変わっていない実行は、両方とも0で残り、遅れて届いた変更が反映された実行は、更新件数が上がります。この2つの数字を見るだけで、パイプラインが正常なのか、ソースに異常があるのかを判断できます。
そして、リトライのコードの中には、データベースの操作だけを入れる必要があります。メール送信や外部API呼び出しが混ざっていると、リトライのたびに、その副作用が繰り返されます。外部呼び出しが必要なら、そちらも冪等性キーを受け取る方式である必要があります。
次のラボですること
自然キーに主キーをかけ、重複した伝票を取り除き、内容ハッシュで変更の有無を判断するupsertを作ります。そして、何も変わっていない再実行と、1件だけ変わった再実行が、それぞれどのように記録されるかを確認します。