TT Lab
はじめる
学ぶ 学習パス コース

データパイプライン

冪等性 — パイプラインは必ずまた回ることになる

TT Labで続きを見る

一言でいうと

パイプラインは失敗し、失敗したらもう一度実行し、そのとき結果が変わるとデータが汚染されるため、再実行の安全性は、選択ではなく設計の前提です。

なぜ必要なのか

バッチが明け方に失敗したとします。朝にもう一度実行します。ところが、失敗した地点がロードの途中だったなら、一部はすでに入っています。そのままもう一度入れると重複が生じ、丸ごと消してからもう一度入れると、その間に入ってきた別のデータまで消えます。

この状況を、毎回人に判断させると、いつか間違いが起きます。そもそも何回入れても同じ結果になるように作るほうがよいです。

どう動くのか

冪等性の出発点は、自然キーです。各行を一意に識別する値があって初めて、「すでに入ってきたもの」を見分けられます。伝票番号、注文番号、イベント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件だけ変わった再実行が、それぞれどのように記録されるかを確認します。