結論から書きます。 Webhook の冪等な受信、Transactional Outbox、冪等キーを正しく実装しても、自社の DB と Stripe のような外部サービスの状態は、いずれズレます。外部が送らなかったイベント、Webhook の再送期間を超える障害、バグ、管理画面からの手動操作が原因です。Reconciliation(整合性チェック) は、そのズレを定期的に検出し、分類し、修復する仕組みです。設計の要点は4つです。状態を3種類に分けて比べること、全件を一定期間内に必ず見る保証を測ること、確認できなかった状態を健全とみなさないこと、そして修復を Reconciler 自身が直接書かず、Webhook と同じ正規の書き込み経路にコマンドとして渡すことです。
Reconciliation は、トランザクションや冪等性の代わりではありません。原子性で防げるズレは原子性で防ぎ、それでも残るズレを後から収束させる最後の安全網です。
コードは PostgreSQL 18.6 / psycopg 3.3 で、「LIMIT 200 が古い行を永遠に見ない」ことの再現を含む12のテストを実行して検証済みです(2026年10月時点)。
コードの前提:psycopg のコードは
psycopg.connect(dsn, autocommit=True)の接続を前提にし、トランザクションの境界をwith conn.transaction():だけで表しています。モジュールの先頭にfrom __future__ import annotationsを置いてください(run_batchの注釈にあるRateLimiterは、9節で後から定義されるため)。
1. どんなズレが起きるのか
具体例を2つ挙げます。
例1:自社では入金待ち、Stripe では入金済み。 payment_intent.succeeded の Webhook を受け取る前にエンドポイントが長時間落ち、Stripe の再送期間(本番で最大3日間)を過ぎた。あるいは、Webhook の処理に「先に処理済みを記録する」バグがあり、イベントが失われた(Webhookの冪等性をTransactional Inboxで実装するの Case A)。顧客は支払ったのに、注文は「入金待ち」のまま止まっている。
例2:自社では入金済み、Stripe では支払いが取り消されている。 自社は paid として商品を出荷したが、その後 Stripe 上で PaymentIntent がキャンセルされていた。原因は、イベントの取りこぼし、管理画面での手動操作、別システムからの API 呼び出しなど様々です。
例1は「自社を外部に追いつかせれば直る」ズレです。例2はお金と出荷が絡む判断を伴い、機械的に paid を取り消せば済む話ではありません。Reconciliation の設計は、この2つを区別するところから始まります。
2. 3つの状態を分ける
ズレを正確に扱うには、「状態」を1つの列で表さず、3つに分けます。
| 状態 | 意味 | 本記事の例での置き場所 |
|---|---|---|
| Desired(意図する状態) | 自社の業務が「こうあるべき」と決めた状態 | payments.status(pending / paid / failed) |
| Known Applied(最後に確認できた外部の状態) | 自社が最後に外部で確認できた状態と、その時刻 | payments.observed_external_status / last_reconciled_at |
| External Authoritative(外部の今の状態) | 外部が今持っている、正とする状態 | Stripe API の PaymentIntent の status |
Known Applied を持つ理由は2つあります。1つは監視です。last_reconciled_at が古い行は「しばらく誰も外部と突き合わせていない行」なので、Reconciler が本当に全件を見ているかを測れます(5章)。もう1つは変化の検出です。前回確認したときと外部の状態が変わったかどうかで、調査の優先度を変えられます。
どれを「正」とするかは、データの種類ごとに決めます。決済の成否は Stripe が正です。一方、「この顧客に割引を適用すべきか」は自社の業務判断が正で、Stripe はその反映先です。前者は外部に合わせて自社を直し、後者は自社に合わせて外部を直します。どちらが正かを決めずに「ズレたら直す」と書くと、2つのシステムが互いを上書きし合います。
3. Reconciler は何をして、何をしないか
Reconciliation を実装するとき、両極端な失敗があります。
- SELECT してログを出すだけ:ズレは見つかるが、誰も見ないログに埋もれ、直らない。
- 何でも直接修復する第二の業務実装:Reconciler が
UPDATE payments SET status = 'paid'を直接書く。Webhook ハンドラと同じ業務ロジック(元帳への計上、通知、在庫の確定)を Reconciler にも書くことになり、2つの実装が少しずつ食い違っていく。Webhook と同時に走ったときに二重計上も起こす。
筆者が推奨する流れは次の通りです。
detect(外部と比較)
↓
classify(ズレの種類と自動修復の可否を決める。純粋関数)
↓
record finding(ズレを記録。同じズレは重複登録しない)
↓
repair command(修復コマンドを Outbox に積む。自動修復できるものだけ)
↓
canonical writer(Webhook と同じ正規の書き込み経路で適用。冪等な状態遷移)
Reconciler の責任は「ズレを見つけて、修復を依頼する」ところまでです。実際の状態変更は、Webhook ハンドラも使う1つの正規の書き込み経路だけが行います。 これで業務ロジックは1か所に保たれ、Webhook と修復が同時に走っても状態遷移の冪等性が守ります。
4. ❌「毎時 LIMIT 200」はなぜ全件を見ないのか
最もよく見る実装です。
# ❌ 壊れた例:「毎時 LIMIT 200」は全件を見る保証にならない。コピーしないでください
def scan_broken(conn: psycopg.Connection, gateway: PaymentGateway, *, batch: int = 200) -> int:
rows = conn.execute(
"""
SELECT id, stripe_payment_intent_id FROM payments
ORDER BY updated_at DESC
LIMIT %s
""",
(batch,),
).fetchall()
for payment_id, pi_id in rows:
gateway.payment_intent_status(pi_id)
conn.execute("UPDATE payments SET last_reconciled_at = now() WHERE id = %s", (payment_id,))
return len(rows)
これを毎時実行すれば、そのうち全件を見る——と思いがちですが、毎回ほぼ同じ200行を見続けます。 ORDER BY updated_at DESC LIMIT 200 が選ぶのは「最近更新された200行」で、古い行は新しい行が増えるほど順位が下がり、二度と選ばれません。
検証では、500行の決済に対してこの関数を10回(10時間分)実行すると、300行は一度も確認されませんでした。ORDER BY が無い場合も同じで、PostgreSQL は ORDER BY の無い LIMIT の結果の順序を保証しないので、どの行が選ばれるかは実行計画次第です。「たまたまいつも同じ行」になることも珍しくありません。
しかも、取りこぼしが起きやすいのはまさに古い行です。Webhook の再送期間を過ぎた古い取引ほど、Reconciliation 以外に救う手段がありません。
4.1 全件を見るのに何周かかるか
全件を見る保証は、件数で考えます。
N = 対象の行数、B = 1回のバッチサイズ、R = 1時間あたりの実行回数、G = 1時間あたりの新規行数
1周に必要な時間 ≈ N ÷ (B × R) ただし B × R > G でなければ、永遠に1周しない
例えば N = 120,000 行、B = 200、R = 12(5分ごと)なら、1周は 120,000 ÷ 2,400 = 50 時間です。「ズレは最大でも約2日以内に見つかる」と言えるのは、この計算が成り立ち、かつ実際に1周が完了していることを測っているときだけです。
5. keyset カーソルで前進し、1周の完了を記録する
5.1 スキーマ
CREATE TABLE payments (
id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
order_id text NOT NULL,
stripe_payment_intent_id text NOT NULL UNIQUE,
status text NOT NULL CHECK (status IN ('pending', 'paid', 'failed')),
version bigint NOT NULL DEFAULT 0,
created_at timestamptz NOT NULL DEFAULT now(),
updated_at timestamptz NOT NULL DEFAULT now(),
observed_external_status text, -- Known Applied: 最後に Stripe で確認できた状態
last_reconciled_at timestamptz -- 最後に外部と突き合わせられた時刻(確認できた時だけ進む)
);
CREATE TABLE reconciliation_cursor (
job text PRIMARY KEY,
last_id bigint NOT NULL DEFAULT 0,
cycle_started_at timestamptz NOT NULL DEFAULT now(),
last_cycle_completed_at timestamptz
);
CREATE TABLE reconciliation_findings (
id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
payment_id bigint NOT NULL REFERENCES payments (id),
kind text NOT NULL,
local_status text NOT NULL,
external_status text,
detected_at timestamptz NOT NULL DEFAULT now(),
resolved_at timestamptz
);
-- 同じズレを毎回のスキャンで重複登録しない(未解決の間は1件)
CREATE UNIQUE INDEX reconciliation_findings_open
ON reconciliation_findings (payment_id, kind) WHERE resolved_at IS NULL;
カーソルは主キー id で前進させます。WHERE id > last_id ORDER BY id LIMIT B は主キーのインデックスを範囲で読むので、テーブルが大きくなっても1回あたりのコストは一定です。OFFSET を使うと、読み飛ばす行数に比例して遅くなるうえ、途中で行が増減すると読み飛ばしや二重読みが起きます。
id が単調増加でない(UUID の主キーなど)場合は、(created_at, id) のような一意で順序のある組をカーソルにします。PostgreSQL の行コンストラクタ比較 WHERE (created_at, id) > (%s, %s) ORDER BY created_at, id が使えます。
5.2 分類:I/O を持たない純粋関数にする
from dataclasses import dataclass
from enum import Enum
class Verdict(Enum):
IN_SYNC = "in_sync"
IN_FLIGHT = "in_flight" # 直近に更新された / 外部が処理中。今は判定しない
REPAIRABLE = "repairable" # 正規の書き込み経路へ修復コマンドを流してよい
NEEDS_REVIEW = "needs_review" # 自動で直さない(お金・出荷が絡む方向のズレ)
UNVERIFIABLE = "unverifiable" # 外部を確認できなかった。健全とはみなさない
# 外部(Stripe PaymentIntent)の状態のうち、まだ遷移中のもの
_EXTERNAL_IN_FLIGHT = {"processing", "requires_action", "requires_confirmation", "requires_capture"}
@dataclass(frozen=True)
class Decision:
verdict: Verdict
kind: str | None = None # finding の種類
repair_command: str | None = None # Outbox に積むコマンド名
def classify(local_status: str, external_status: str) -> Decision:
if local_status == "pending" and external_status in _EXTERNAL_IN_FLIGHT:
# 自社も外部もまだ確定していない。自社が paid なのに外部が遷移中、は下の ("paid", _) へ
return Decision(Verdict.IN_FLIGHT)
match (local_status, external_status):
case ("paid", "succeeded") | ("failed", "canceled") | ("pending", "requires_payment_method"):
return Decision(Verdict.IN_SYNC)
case ("pending", "succeeded"):
# Webhook の取りこぼし。Webhook と同じ正規の書き込み経路で「入金済み」に進める
return Decision(Verdict.REPAIRABLE, "missed_success", "ApplyPaymentSucceeded")
case ("pending", "canceled"):
return Decision(Verdict.REPAIRABLE, "missed_cancel", "ApplyPaymentCanceled")
case ("paid", _):
# 自社は入金済みとして扱った(出荷・付与した可能性)が、外部では成功していない
return Decision(Verdict.NEEDS_REVIEW, "paid_locally_not_externally")
case _:
return Decision(Verdict.NEEDS_REVIEW, f"unexpected:{local_status}->{external_status}")
分類の方針:
- 自社を外部に追いつかせる方向(
pending→ 外部ではsucceeded/canceled)はREPAIRABLE。Webhook が届いていれば起きていたはずの状態遷移を、同じ経路で起こすだけだからです。 - 自社が先に進んでいる方向(自社
paid、外部は成功していない)はNEEDS_REVIEW。出荷の停止、顧客への連絡、返金の要否は業務の判断で、Reconciler が決めてよいことではありません。 - 自社が
pendingで、外部が遷移中(processingなど)はIN_FLIGHT。今判定すると、数秒後に正しくなる状態をズレとして報告してしまいます。自社がpaidなのに外部が遷移中(認証待ちのまま放置されたrequires_actionなど)はIN_FLIGHTにせずNEEDS_REVIEWにします。IN_FLIGHTにすると外部の状態は確認できているので最終確認時刻は進み続け、「遅延の監視は正常、実は入金されていない」という偽の健全が永遠に続くからです(検証済み)。 - 想定外の組み合わせは
NEEDS_REVIEWに倒します。未知のものを「健全」に分類しないためです。
分類を純粋関数にしておくと、全ての組み合わせを I/O 無しの表形式テストで網羅できます。PaymentIntent の status の値の一覧は、Stripe API リファレンスで確認してください。
5.3 1バッチの実行
from dataclasses import field
from datetime import timedelta
from typing import Protocol
import psycopg
from psycopg.types.json import Jsonb
class GatewayUnavailable(Exception):
"""タイムアウト・429・5xx。外部状態を確認できなかった。"""
class ExternalNotFound(Exception):
pass
class PaymentGateway(Protocol):
"""Stripe SDK を包む境界。SDK の例外を上の2つに変換するのはこの実装の責任。"""
def payment_intent_status(self, payment_intent_id: str) -> str: ...
@dataclass
class BatchReport:
scanned: int = 0
counts: dict[Verdict, int] = field(default_factory=lambda: {v: 0 for v in Verdict})
wrapped: bool = False
class ConcurrentRun(Exception):
"""別の reconciler がカーソルを先に進めた。このバッチの書き込みは捨てる。"""
def run_batch(
conn: psycopg.Connection,
gateway: PaymentGateway,
*,
job: str = "payments-vs-stripe",
batch: int = 200,
grace: timedelta = timedelta(minutes=10),
limiter: RateLimiter | None = None,
) -> BatchReport:
report = BatchReport()
# 1) 短いトランザクションでカーソル位置と対象行を読む(外部 API 中にロックを握らない)
with conn.transaction():
conn.execute("INSERT INTO reconciliation_cursor (job) VALUES (%s) ON CONFLICT DO NOTHING", (job,))
cursor_row = conn.execute(
"SELECT last_id FROM reconciliation_cursor WHERE job = %s", (job,)
).fetchone()
assert cursor_row is not None # 直前の INSERT ... ON CONFLICT で必ず存在する
cursor_id: int = cursor_row[0]
rows = conn.execute(
"""
SELECT id, stripe_payment_intent_id, status,
updated_at > now() - %s AS recently_changed
FROM payments
WHERE id > %s
ORDER BY id
LIMIT %s
""",
(grace, cursor_id, batch),
).fetchall()
if not rows:
# 末尾まで到達 = 1周完了。先頭へ巻き戻す(CAS で他の実行と競合しない)
with conn.transaction():
cur = conn.execute(
"""
UPDATE reconciliation_cursor
SET last_id = 0, last_cycle_completed_at = now(), cycle_started_at = now()
WHERE job = %s AND last_id = %s
""",
(job, cursor_id),
)
if cur.rowcount != 1:
raise ConcurrentRun(job)
report.wrapped = True
return report
# 2) トランザクションの外で外部の状態を読む
results: list[tuple[int, str, Decision, str | None]] = []
for payment_id, pi_id, local_status, recently_changed in rows:
if recently_changed:
results.append((payment_id, local_status, Decision(Verdict.IN_FLIGHT), None))
continue
if limiter:
limiter.wait()
try:
external = gateway.payment_intent_status(pi_id)
except GatewayUnavailable:
results.append((payment_id, local_status, Decision(Verdict.UNVERIFIABLE), None))
continue
except ExternalNotFound:
results.append(
(payment_id, local_status, Decision(Verdict.NEEDS_REVIEW, "missing_externally"), None)
)
continue
results.append((payment_id, local_status, classify(local_status, external), external))
# 3) 結果の記録・修復コマンド・カーソル前進を1つのトランザクションで
with conn.transaction():
for payment_id, local_status, decision, observed in results:
report.counts[decision.verdict] += 1
if observed is not None:
# 外部を実際に確認できた行だけ「検証済み」を進める(UNVERIFIABLE は進めない)
conn.execute(
"""
UPDATE payments SET observed_external_status = %s, last_reconciled_at = now()
WHERE id = %s
""",
(observed, payment_id),
)
if decision.kind is None:
continue
finding = conn.execute(
"""
INSERT INTO reconciliation_findings (payment_id, kind, local_status, external_status)
VALUES (%s, %s, %s, %s)
ON CONFLICT (payment_id, kind) WHERE resolved_at IS NULL DO NOTHING
RETURNING id
""",
(payment_id, decision.kind, local_status, observed),
).fetchone()
if finding is not None and decision.repair_command is not None:
# 直接 UPDATE しない。Webhook と同じ「正規の書き込み経路」へコマンドとして渡す
conn.execute(
"""
INSERT INTO outbox (aggregate_type, aggregate_id, event_type, payload)
VALUES ('Payment', %s, %s, %s)
""",
(
str(payment_id),
decision.repair_command,
Jsonb({"payment_id": payment_id, "finding_id": finding[0],
"expected_local_status": local_status}),
),
)
cur = conn.execute(
"UPDATE reconciliation_cursor SET last_id = %s WHERE job = %s AND last_id = %s",
(rows[-1][0], job, cursor_id),
)
if cur.rowcount != 1:
raise ConcurrentRun(job) # 例外でトランザクションごと捨てる
report.scanned = len(rows)
return report
このコードの設計判断:
- 外部 API はトランザクションの外で呼ぶ。 200件 × 数百ミリ秒の API 呼び出しの間、DB の接続もロックも握らない。
- 結果の記録とカーソルの前進を1つのトランザクションにする。 途中で落ちれば、カーソルも記録も進まず、次回同じバッチをやり直す。finding は部分一意インデックスで重複しないので、やり直しても安全。
- カーソルの前進を CAS(Compare-And-Set)にする。
WHERE last_id = <読んだときの値>で更新し、0行なら別の Reconciler が先に進めたと判断してバッチごと捨てる。重複起動しても、記録やコマンドが二重にならない(検証済み)。 - 猶予期間(
grace)内に更新された行は判定しない。 直前に Webhook で更新された行は、外部との反映の時間差でズレに見えることがある。その行は次の周で見る。 - 修復コマンドは、finding を新規に作ったときだけ積む。 同じズレが未解決のまま次の周で見つかっても、コマンドは1つのまま(検証済み)。
5.4 1周の保証を測る
カーソルが末尾に達すると last_cycle_completed_at が記録され、先頭に戻ります。検証では、450行をバッチ200で処理すると、200 → 200 → 50 → 0(1周完了・先頭へ)と進み、全行の last_reconciled_at が埋まりました。
ただし、カーソルが1周したことは「全行を確認できた」ことを意味しません。外部 API が落ちていた行は UNVERIFIABLE として飛ばされているからです。だから、本当の保証は行ごとの最終確認時刻で測ります。
-- reconciliation lag:最も長く外部と突き合わせていない行の古さ(最重要)
SELECT max(now() - coalesce(last_reconciled_at, created_at)) AS reconciliation_lag
FROM payments;
-- 1周の所要時間と、前回の完了からの経過時間
SELECT job,
last_cycle_completed_at,
now() - last_cycle_completed_at AS since_last_cycle,
now() - cycle_started_at AS current_cycle_age
FROM reconciliation_cursor;
reconciliation_lag に SLO(例:72時間)を置き、超えたらアラートにします。これで「Reconciler は動いているが、外部 API のエラーで実は何も確認できていない」状態も検出できます。
6. 「確認できない」を「健全」と数えない
外部 API がタイムアウトしたとき、try/except でエラーを握りつぶして次の行に進むと、その行は何も起きなかった=ズレていないように見えます。Stripe が1日落ちていれば、Reconciler は1日中「ズレ0件」を報告します。これが false healthy(偽の健全) です。
本記事のコードでは、GatewayUnavailable の行を UNVERIFIABLE として数え、last_reconciled_at を進めません。検証でも、外部が全件エラーを返したとき、UNVERIFIABLE が10件、IN_SYNC が0件となり、どの行の最終確認時刻も進みませんでした。
監視では、バッチごとの UNVERIFIABLE の割合をメトリクスとして出し、IN_SYNC の件数だけで健全性を判断しないようにします。
外部で見つからない(404)行は、UNVERIFIABLE ではなく NEEDS_REVIEW(missing_externally)にしています。「確認できた結果、存在しなかった」は明確な事実で、テストデータの混入、別アカウントのキー、削除など、人が原因を調べるべき状態だからです。
7. 修復:正規の書き込み経路に渡す
Reconciler が Outbox に積んだ ApplyPaymentSucceeded は、Dispatcher が取り出し、Webhook ハンドラと同じ関数で適用します。
def apply_payment_succeeded(conn: psycopg.Connection, payment_id: int, *, finding_id: int | None = None) -> bool:
"""pending → paid の状態遷移。すでに進んでいれば何もしない(Webhook と修復の競合に安全)。
修復コマンドから呼ばれたときは、遷移が no-op でも finding を解決済みにする
(目的の状態に達していることに変わりはないため)。
"""
with conn.transaction():
cur = conn.execute(
"""
UPDATE payments SET status = 'paid', version = version + 1, updated_at = now()
WHERE id = %s AND status = 'pending'
""",
(payment_id,),
)
if finding_id is not None:
conn.execute(
"UPDATE reconciliation_findings SET resolved_at = now() WHERE id = %s AND resolved_at IS NULL",
(finding_id,),
)
return cur.rowcount == 1
WHERE status = 'pending' が、修復と Webhook の競合を安全にしています。Reconciler がズレを見つけてからコマンドが適用されるまでの間に、遅れて届いた Webhook が先に paid にしていれば、修復コマンドは0行更新の no-op になります(検証済み)。実際の業務では、この関数の中で元帳への計上や通知の Outbox への記録も行います。Webhook と修復が同じ関数を通るので、その処理は1か所にしか書かれません。
修復コマンドを Outbox 経由にする理由は、Reconciler のトランザクション(finding の記録)と、修復の依頼を原子的にするためです。Reconciler が直接 Dispatcher のキューやメッセージブローカーに送ると、「finding は記録したがコマンドは送れなかった」という2つの書き込みの問題がまた生まれます(Transactional Outbox)。Outbox を複数の worker で配送する方法は Outbox Dispatcher の lease と fencing token で扱っています。
修復コマンドから呼ぶときは finding_id を渡し、遷移が no-op でも finding を解決済みにします。finding が未解決のまま残ると、部分一意インデックスのせいで同じズレが二度と登録されず、修復コマンドが dead になった場合に永久に再試行されなくなるからです。運用では、「一定時間(例:24時間)以上未解決の finding」をアラートにします。
NEEDS_REVIEW の finding は、管理画面やチケットで人に渡します。人が対応した結果も、正規の書き込み経路を通して適用し、finding の resolved_at を埋めます。誰がいつどう判断したかを残しておくことが、決済の監査では重要です。
8. 反対方向のスキャン:外部にだけ存在するリソース
ここまでのスキャンは「自社の行を起点に外部を見る」方向です。この方向では、外部にだけ存在するリソース(orphan)を見つけられません。例えば、自社のトランザクションが ROLLBACK したのに、その前に Stripe 側で PaymentIntent を作ってしまっていた場合です。
これを見つけるには、外部の一覧を起点に自社を照合する、反対方向のスキャンを別のジョブとして持ちます。Stripe の一覧 API は starting_after(オブジェクトID)によるカーソル型のページネーションで、新しい順に最大100件ずつ返します。1回のジョブで扱う期間(例:作成日時が直近7日)を区切り、ページのカーソルを永続化して、同じく少しずつ前進させます。
Stripe の Search API で metadata を検索する方法もありますが、公式ドキュメントは検索結果への反映が通常1分未満で、障害時は遅れ得るとしており、作成直後の確認(read-after-write)に使わないよう注意しています。Reconciliation のように時間をおいて照合する用途には使えますが、「見つからない=存在しない」と即断しないようにします。
9. レート制限:本番の API 予算を食い尽くさない
Reconciler は外部 API を大量に呼びます。Stripe の公式ドキュメントでは、レート制限は Stripe アカウント単位で、本番の全体上限は毎秒100リクエスト、個別のエンドポイントは特記が無い限り毎秒25リクエストとされています(2026年10月時点。上限は変わり得るので公式ドキュメントで確認してください)。
この上限は、決済そのものの API 呼び出しと共有です。 Reconciler が上限まで使えば、顧客の決済が 429 で失敗します。Reconciler には上限のごく一部(例:毎秒5〜10リクエスト)だけを割り当てます。
import time
class RateLimiter:
"""単純な固定間隔リミッタ。本番 API のレート上限の一部だけを使う。"""
def __init__(self, rps: float) -> None:
self._interval = 1.0 / rps
self._next = time.monotonic()
def wait(self) -> None:
now = time.monotonic()
if now < self._next:
time.sleep(self._next - now)
self._next = max(now, self._next) + self._interval
この割り当ては、4.1節の計算に直結します。毎秒5リクエストなら1時間に18,000件で、12万行の1周は約7時間です。行数が増えて SLO を満たせなくなったら、レートを上げるのではなく、突き合わせる対象を絞る(例:直近90日の取引と、未完了の状態の行だけを高頻度に、それ以外は低頻度に)ことを先に検討します。
429 を受け取ったら GatewayUnavailable として扱い、そのバッチの残りの呼び出しは止めて次回に回すのが安全です(本記事のコードは単純化のため行ごとに判定しています)。
10. テスト
PostgreSQL 18 のコンテナに対して実行し、すべて通過を確認したテストです。外部 API は、状態を辞書で返すフェイクに差し替えています。
| テスト | 期待する結果 |
|---|---|
| 分類の表 | pending×succeeded → REPAIRABLE、paid×canceled → NEEDS_REVIEW、processing → IN_FLIGHT |
❌ ORDER BY updated_at DESC LIMIT 200 を10回(500行) | 300行が一度も確認されない(失敗の再現) |
| keyset カーソル(450行・バッチ200) | 200 → 200 → 50 → 0(1周完了)、全行の最終確認時刻が埋まる |
| 外部 API が全件エラー | UNVERIFIABLE 10件、IN_SYNC 0件、最終確認時刻は進まない |
| ズレの検出を2周 | finding とコマンドはズレ1種につき1件のまま |
自社 paid・外部 canceled | NEEDS_REVIEW の finding のみ。修復コマンドは積まない |
| 猶予期間内に更新された行 | IN_FLIGHT、外部 API は呼ばない |
| 別の Reconciler が先にカーソルを進めた | ConcurrentRun でバッチを捨てる |
| 修復と Webhook の競合 | 先に適用した側だけが状態を変え、もう一方は no-op |
自社 paid・外部 requires_action | IN_FLIGHT ではなく NEEDS_REVIEW |
| 修復コマンドの適用 | finding が解決済みになり、同じズレを再び検出できる |
| レートリミッタ | 毎秒50リクエストで11回の呼び出しに0.19秒以上かかる |
def test_broken_scan_starves_old_rows(conn):
seed(conn, 500)
gw = FakeStripe({f"pi_{i}": "succeeded" for i in range(1, 501)})
for _ in range(10): # 10時間分
scan_broken(conn, gw)
never = conn.execute("SELECT count(*) FROM payments WHERE last_reconciled_at IS NULL").fetchone()[0]
assert never == 300 # 10回回しても 300 行は一度も見られない
11. 使い方を誤らないために
- Reconciliation をトランザクションの代わりにしない。 「どうせ後で直るから」と原子性を省くと、ズレの量が Reconciler の処理能力を超え、
NEEDS_REVIEWが人の手に負えなくなる。 - Reconciler に業務ロジックを書かない。 修復は正規の書き込み経路へのコマンドにする。
- どちらが正かを、データの種類ごとに決めてから書く。 決めずに双方向で直すと、互いを上書きし合う。
- 自動修復は「自社を外部に追いつかせる方向」から始める。 お金・出荷・権限が絡む逆方向は、人の確認を経る。
- 件数ではなく時間で保証を測る。 「1回に200件見る」ではなく「全件が72時間以内に確認されている」を監視する。
まとめ
- Reconciliation は、原子性と冪等性で防ぎきれないズレを後から収束させる安全網で、トランザクションの代替ではない。
- 状態を Desired / Known Applied / External Authoritative の3つに分け、データの種類ごとにどれが正かを決める。
- 「毎時 LIMIT 200」は同じ行を見続ける。keyset カーソルで前進し、1周の完了と行ごとの最終確認時刻で保証を測る。
- 外部 API の失敗は「確認できない」であって「健全」ではない。最終確認時刻を進めず、割合を監視する。
- 修復は Outbox 経由で正規の書き込み経路に渡す。お金や出荷が絡む方向は人の確認に回す。
- 外部の API 予算は本番の決済と共有。Reconciler には上限の一部だけを割り当てる。
外部 API の呼び出しが「成功したかどうか分からない」状態を、冪等キーと Reconciliation で1回の効果に収束させる考え方は、Exactly-onceは何を保証するのかで扱います。Stripe の Webhook とサブスクリプションの状態機械は StripeのWebhookと冪等性を本番品質で実装する にまとめています。