結論から書きます。 Webhookの冪等性は「イベントIDにUNIQUE制約を張る」だけでは完成しません。受信の記録と業務の更新を、同じデータベーストランザクションでコミットする必要があります。先に「処理済み」を記録すると業務処理の失敗でイベントが永久に失われ、後で記録するとクラッシュ時に二重処理になるからです。この「同じトランザクション」に入れる受信側の設計を Transactional Inbox と呼びます。そのうえで、別のイベントとして届く同じ業務事実や順不同に備えて、ハンドラ自体も冪等にします。
この記事では、2つの失敗パターンを実際に再現し、PostgreSQL の INSERT ... ON CONFLICT DO NOTHING を使った Transactional Inbox を、同期版と非同期版の両方で実装します。コードは PostgreSQL 18.6 / psycopg 3.3 / stripe-python 16 で、同時重複配信・ワーカーのクラッシュを含む18のテストを実行して検証済みです(2026年10月時点)。DynamoDB による重複排除との比較も、どちらかを「間違い」とせずに整理します。
前提:Stripe の仕様は公式ドキュメントの2026年10月時点の記述に基づきます。Stripe 以外のプロバイダの再送・順序の仕様は、それぞれの公式ドキュメントで確認してください。Stripe の仕様を Webhook 一般の仕様として扱わないよう注意して書きます。
1. 前提:Webhookは「一回だけ」届くとは限らない
Stripe の公式ドキュメントは、配信について次のことを明記しています(2026年10月時点)。
- 再送:本番環境では、配信に失敗したイベントを最大3日間、指数バックオフで再送する。
- 重複:Webhook エンドポイントは同じイベントを複数回受信することがある。処理済みのイベントIDを記録し、記録済みのものは処理しないことで防ぐよう推奨している。
- 別オブジェクトとしての重複:場合によっては2つの別々の Event オブジェクトが生成されて送られる。その重複は
data.objectのIDとevent.typeで識別するよう案内している。 - 順序:イベントが生成された順に配信されることは保証しない。
再送そのものも避けられません。Stripe は 2xx を受け取れなければ失敗とみなします。ハンドラは正しく処理したのに、レスポンスがタイムアウトやネットワーク断で届かなければ、同じイベントがもう一度届きます。つまり受信側は、同じイベントが複数回、しかも同時に届く可能性を前提に作る必要があります。
2. 2つの失敗パターン
「イベントIDを記録して重複を弾く」という方針自体は正しいのですが、記録と業務処理をどの順序で、どのトランザクションでコミットするかを間違えると、正反対の2つの壊れ方をします。
題材は invoice.paid です。入金を元帳(ledger_entries)に計上し、請求書を paid にします。
Case A:先に「処理済み」をコミット → イベントの永久消失
# ❌ 壊れた例:受信記録を先に別トランザクションでコミットしている
def receive_mark_first_broken(conn, *, endpoint, event):
with conn.transaction():
inserted = conn.execute(
"INSERT INTO webhook_inbox (provider, endpoint, event_id, event_type, payload) "
"VALUES ('stripe', %s, %s, %s, %s) ON CONFLICT DO NOTHING RETURNING event_id",
(endpoint, event.id, event.type, Jsonb(event.payload)),
).fetchone()
if inserted is None:
return Outcome.DUPLICATE # 記録済み=処理済みとみなして 200
with conn.transaction():
apply_event(conn, event) # ← ここで失敗すると…
return Outcome.PROCESSED
1回目の配信: 受信記録 COMMIT → 業務処理で DB エラー → ROLLBACK → 500 を返す
2回目の配信: 受信記録の INSERT が衝突 → 「重複」と判断 → 200 を返す
→ 入金は一度も計上されず、Stripe はもう再送しない
これは at-most-once(高々1回)側の失敗です。重複は起きない代わりに、取りこぼしが起きます。検証でも、1回目の業務処理を失敗させると、2回目は DUPLICATE を返し、元帳は0件のままでした。決済では、二重課金と同じくらい深刻な障害です。
Case B:業務処理をコミットしてから記録 → 二重処理
# ❌ 壊れた例:業務処理のコミットと受信記録が別トランザクション
def receive_mark_after_broken(conn, *, endpoint, event):
seen = conn.execute(
"SELECT 1 FROM webhook_inbox WHERE provider='stripe' AND endpoint=%s AND event_id=%s",
(endpoint, event.id),
).fetchone()
if seen:
return Outcome.DUPLICATE
with conn.transaction():
apply_event(conn, event) # 業務処理を COMMIT
# ← ここでプロセスが落ちる(デプロイ、OOM、タイムアウト)
with conn.transaction():
conn.execute("INSERT INTO webhook_inbox ...")
return Outcome.PROCESSED
1回目の配信: 業務処理 COMMIT(入金を計上)→ 受信記録の前にクラッシュ → 2xx は返らない
2回目の配信: 受信記録が無い → 未処理と判断 → 業務処理をもう一度 COMMIT
→ 入金が2回計上される
これは at-least-once(少なくとも1回)側の失敗です。検証でも、受信記録の直前で処理を中断させてから再送すると、元帳は2件になりました。SELECT してから INSERT しているので、同時に2本届いた場合も両方が「未処理」と判断して二重処理になります(TOCTOU)。
2つのパターンの原因は同じです。 受信記録と業務更新という2つの書き込みを、別々にコミットしていることです。どちらを先にしても、間で落ちれば片方だけが残ります。
3. 解:受信記録と業務更新を1つのトランザクションに入れる(Transactional Inbox)
受信記録の INSERT と業務更新を、同じトランザクションでコミットします。
BEGIN
INSERT INTO webhook_inbox (event_id ...) ON CONFLICT DO NOTHING RETURNING ...
→ 行が返らなければ重複。何もせず COMMIT して 200
UPDATE invoices / INSERT ledger_entries(業務更新)
INSERT INTO outbox(下流への通知が要るなら)
UPDATE webhook_inbox SET processed_at = now()
COMMIT → 200
- 業務処理が失敗すれば、受信記録も一緒に ROLLBACK される。受信記録が残らないので、再送は未処理として正しく再処理される(Case A が起きない)。
- COMMIT が成功すれば、業務更新と受信記録が両方残る。その後プロセスが落ちて 2xx が返らなくても、再送は重複として弾かれる(Case B が起きない)。
これが Transactional Inbox です。送信側で「業務更新とイベントの記録を同じトランザクションに入れる」Transactional Outbox と鏡像の関係にあります。
| Transactional Outbox(送信側) | Transactional Inbox(受信側) | |
|---|---|---|
| 同じトランザクションに入れるもの | 業務更新 + 送るべきイベント | 受け取ったイベントID + 業務更新 |
| 防ぐ失敗 | 業務は更新したのにイベントが出ない/出たのに業務は無い | 受け取ったのに処理されない/2回処理される |
| 残る性質 | 発行は at-least-once(重複発行はあり得る) | 重複配信を吸収し、業務効果を1回に収束させる |
なお、Inbox はデータベースが1つの場合にだけこの原子性を持ちます。受信記録と業務データが別のストアにあると、また2つの書き込みに戻ります(7章の DynamoDB との比較で扱います)。
4. スキーマ
-- 受信台帳。主キーが「同じ配信は1回だけ受け付ける」を DB に裁かせる
CREATE TABLE webhook_inbox (
provider text NOT NULL, -- 'stripe' など。プロバイダ間の ID 衝突を分離
endpoint text NOT NULL, -- 同じイベントを受ける別エンドポイントを分離
event_id text NOT NULL, -- Stripe なら evt_...
event_type text NOT NULL,
payload jsonb NOT NULL, -- 署名検証済みの生 JSON
received_at timestamptz NOT NULL DEFAULT now(),
available_at timestamptz NOT NULL DEFAULT now(), -- 非同期処理の再試行時刻
attempts integer NOT NULL DEFAULT 0,
last_error text,
processed_at timestamptz, -- NULL = 未処理
PRIMARY KEY (provider, endpoint, event_id)
);
CREATE INDEX webhook_inbox_unprocessed
ON webhook_inbox (available_at) WHERE processed_at IS NULL;
CREATE TABLE invoices (
id text PRIMARY KEY, -- in_...
customer_id text NOT NULL,
status text NOT NULL CHECK (status IN ('draft', 'open', 'paid', 'void', 'uncollectible')),
amount_paid bigint NOT NULL DEFAULT 0,
paid_at timestamptz,
updated_at timestamptz NOT NULL DEFAULT now()
);
-- 下流への通知(構造は Transactional Outbox の記事と同じ)
CREATE TABLE outbox (
id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
aggregate_type text NOT NULL,
aggregate_id text NOT NULL,
event_type text NOT NULL,
payload jsonb NOT NULL,
created_at timestamptz NOT NULL DEFAULT now(),
published_at timestamptz
);
-- 業務レベルの冪等性:同じ請求の入金を2回計上しない(別 Event オブジェクトで届いても)
CREATE TABLE ledger_entries (
id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
invoice_id text NOT NULL,
kind text NOT NULL,
amount bigint NOT NULL,
UNIQUE (invoice_id, kind)
);
設計判断を3つ説明します。
主キーに endpoint を含める理由:同じ Stripe アカウントに複数のエンドポイント(例:課金処理用と分析用)を登録すると、同じイベントが両方に届きます。両方とも処理すべきなので、イベントIDだけを主キーにすると、片方のエンドポイントがもう片方の処理を「重複」として潰します。検証でも、endpoint を分けると両方とも PROCESSED になります。プロバイダ(provider)も含めるのは、別のサービスのイベントIDが偶然一致しても混ざらないようにするためです。
payload を保存する理由:非同期処理(6章)でワーカーが処理するため、そして障害調査で「何が届いたか」を後から確認するためです。個人情報を含むイベントを保存するなら、保持期間(9章)と閲覧権限を決めておきます。
元帳に UNIQUE (invoice_id, kind) を張る理由:受信台帳が防ぐのは「同じイベントIDの重複」だけです。Stripe が別の Event オブジェクトとして同じ事実を送ってきた場合、イベントIDは違うので受信台帳を素通りします。業務キー(請求書ID)の一意制約が、この重複を止める第二の防衛線になります(5.2節)。
5. 同期版の実装:受信して、その場で処理する
業務更新が DB の中で短時間に終わるなら、これが最も単純で強い形です。
コードの前提:本記事の psycopg のコードは
psycopg.connect(dsn, autocommit=True)の接続を前提にし、トランザクションの境界をwith conn.transaction():だけで表しています。autocommit が無効な接続でtransaction()の外で SQL を実行すると暗黙のトランザクションが始まり、その後のtransaction()は SAVEPOINT になって、ここで説明するコミットのタイミングが成り立ちません。
5.1 署名検証:生のボディで行う
import json
from dataclasses import dataclass
from typing import Any
import stripe
@dataclass(frozen=True)
class VerifiedEvent:
"""署名検証を通過したイベント。これ以外の型は業務ロジックへ入れない。"""
id: str
type: str
payload: dict[str, Any]
class InvalidWebhook(Exception):
pass
def verify_stripe_event(raw_body: bytes, sig_header: str | None, secret: str) -> VerifiedEvent:
"""フレームワークがボディを dict に変換する前の bytes で署名を検証する。"""
try:
stripe.Webhook.construct_event(raw_body, sig_header, secret)
except (ValueError, stripe.SignatureVerificationError) as exc:
raise InvalidWebhook(str(exc)) from exc
payload = json.loads(raw_body)
return VerifiedEvent(id=payload["id"], type=payload["type"], payload=payload)
Stripe の署名は、受信した生のバイト列に対して計算されています。JSON をパースしてから再シリアライズすると、空白やキーの並びが変わって署名が一致しません。検証でも、内容が同じでも区切り文字を変えて再シリアライズしたボディは InvalidWebhook になりました。署名検証は、イベントの内容を信用してよいかを決める唯一の境界です。検証前のデータを DB に書かないでください。
Stripe の SDK は、署名のタイムスタンプと現在時刻の差を既定で5分まで許容します。公式ドキュメントは許容値を 0 にしないよう注意しています(0 は鮮度チェックを無効にする)。サーバーの時計は NTP で同期しておきます。
5.2 ハンドラ自体を冪等・順序非依存に書く
def _apply_invoice_created(conn: psycopg.Connection, inv: dict[str, Any]) -> None:
# 作成は「無ければ入れる」だけ。paid の後に遅れて届いても状態を巻き戻さない
conn.execute(
"""
INSERT INTO invoices (id, customer_id, status)
VALUES (%(id)s, %(customer)s, 'open')
ON CONFLICT (id) DO NOTHING
""",
{"id": inv["id"], "customer": inv["customer"]},
)
def _apply_invoice_paid(conn: psycopg.Connection, inv: dict[str, Any]) -> None:
# paid は created より先に届くことがある → upsert。void 済みの請求は上書きしない
conn.execute(
"""
INSERT INTO invoices (id, customer_id, status, amount_paid, paid_at)
VALUES (%(id)s, %(customer)s, 'paid', %(amount_paid)s, now())
ON CONFLICT (id) DO UPDATE
SET status = 'paid',
amount_paid = EXCLUDED.amount_paid,
paid_at = COALESCE(invoices.paid_at, EXCLUDED.paid_at),
updated_at = now()
WHERE invoices.status NOT IN ('paid', 'void')
""",
{"id": inv["id"], "customer": inv["customer"], "amount_paid": inv["amount_paid"]},
)
row = conn.execute("SELECT status FROM invoices WHERE id = %s", (inv["id"],)).fetchone()
if row is None or row[0] != "paid":
# void 済みの請求への入金は異常系。計上せず、Reconciliation で人の確認に回す
return
booked = conn.execute(
"""
INSERT INTO ledger_entries (invoice_id, kind, amount)
VALUES (%(id)s, 'payment', %(amount_paid)s)
ON CONFLICT (invoice_id, kind) DO NOTHING
RETURNING id
""",
{"id": inv["id"], "amount_paid": inv["amount_paid"]},
).fetchone()
if booked is not None:
# 下流(領収書メール等)への通知は Outbox へ。業務更新と同じトランザクションで書く
conn.execute(
"""
INSERT INTO outbox (aggregate_type, aggregate_id, event_type, payload)
VALUES ('Invoice', %(id)s, 'InvoicePaymentBooked', %(payload)s)
""",
{"id": inv["id"], "payload": Jsonb({"invoice_id": inv["id"], "amount": inv["amount_paid"]})},
)
HANDLERS = {
"invoice.created": _apply_invoice_created,
"invoice.paid": _apply_invoice_paid,
}
ここで効いている設計は4つです。
- 状態遷移として書く:
invoice.createdは「無ければ作る」だけで、既存の状態を変えません。invoice.paidはpaid/void以外のときだけpaidに進めます。だからinvoice.paidがinvoice.createdより先に届いても、最終状態はpaidになります(検証済み)。 createdで順序を判定しない:Stripe の公式ドキュメントは、createdが秒単位で別のイベントが同じ値を持ち得るため、順序や処理済みの判定にcreatedを使わないよう明記しています(2026年10月時点)。順序が本当に必要な判断(サブスクリプションの最新状態など)は、Stripe API から最新のオブジェクトを取得して決めます。公式ドキュメントも、先に届いたイベントの情報から関連オブジェクトを API で取得する方法を案内しています。- 業務キーの一意制約で効果を1回にする:
ledger_entriesのUNIQUE (invoice_id, kind)があるので、別の Event オブジェクト(evt_Aとevt_B)で同じ請求の入金が届いても、計上は1回です(検証済み)。 - 請求書が
paidになったときだけ計上する:void済みの請求書にinvoice.paidが届いた場合、upsert は状態を変えません。そのまま元帳に計上すると「無効な請求への入金」が記録されるので、状態を読み直してpaidでなければ計上しません。この組み合わせは異常系なので、Reconciliation で人の確認に回します(検証済み)。
5.3 受信処理本体
from enum import Enum
class Outcome(Enum):
PROCESSED = "processed"
DUPLICATE = "duplicate"
IGNORED = "ignored"
def apply_event(conn: psycopg.Connection, event: VerifiedEvent) -> bool:
handler = HANDLERS.get(event.type)
if handler is None:
return False
handler(conn, event.payload["data"]["object"])
return True
def receive_sync(
conn: psycopg.Connection, *, provider: str, endpoint: str, event: VerifiedEvent
) -> Outcome:
"""受信記録と業務更新を1つのトランザクションでコミットする。"""
with conn.transaction():
inserted = conn.execute(
"""
INSERT INTO webhook_inbox (provider, endpoint, event_id, event_type, payload)
VALUES (%s, %s, %s, %s, %s)
ON CONFLICT (provider, endpoint, event_id) DO NOTHING
RETURNING event_id
""",
(provider, endpoint, event.id, event.type, Jsonb(event.payload)),
).fetchone()
if inserted is None:
return Outcome.DUPLICATE
handled = apply_event(conn, event)
conn.execute(
"""
UPDATE webhook_inbox SET processed_at = now()
WHERE provider = %s AND endpoint = %s AND event_id = %s
""",
(provider, endpoint, event.id),
)
return Outcome.PROCESSED if handled else Outcome.IGNORED
conn.transaction() のブロック内で例外が起きれば、psycopg はトランザクションを ROLLBACK して例外を再送出します。呼び出し側はそれを 5xx に変換し、Stripe の再送に任せます。
5.4 同時に2本届いたらどうなるか
同じイベントが同時に2本届いた場合の動きが、Transactional Inbox の核心です。
時刻 配信 1(接続1) 配信 2(接続2)
t1 BEGIN; INSERT inbox(evt_1) → 成功(未コミット)
t2 BEGIN; INSERT inbox(evt_1) → 一意インデックスで待機
t3 業務更新
t4 COMMIT
t5 → 衝突が確定。ON CONFLICT DO NOTHING → 0行
t6 COMMIT → DUPLICATE(200)
PostgreSQL の INSERT ... ON CONFLICT は、同じキーを挿入中の未コミットのトランザクションがあると、その結果が確定するまで待ちます。1本目がコミットすれば2本目は衝突として DO NOTHING になり、1本目が ROLLBACK すれば2本目の挿入が成功して処理を引き継ぎます。
検証では、同じイベントを8本の接続から同時に送ると「PROCESSED 1本・DUPLICATE 7本・元帳1件」になりました。1本目の業務処理を 0.5 秒後に失敗させるケースでは、2本目は1本目の ROLLBACK まで待ってから PROCESSED になり、元帳は1件でした。アプリ側でロックを書かなくても、一意インデックスが同時重複を直列化してくれます。
5.5 HTTP の境界(FastAPI の例)
import os
from fastapi import FastAPI, Request, Response
from psycopg_pool import ConnectionPool
from starlette.concurrency import run_in_threadpool
app = FastAPI()
pool = ConnectionPool(os.environ["DATABASE_URL"], kwargs={"autocommit": True}, open=True)
WEBHOOK_SECRET = os.environ["STRIPE_WEBHOOK_SECRET"] # ハードコードしない
def _handle(event: VerifiedEvent) -> Outcome:
with pool.connection() as conn:
return receive_sync(conn, provider="stripe", endpoint="billing", event=event)
@app.post("/webhooks/stripe")
async def stripe_webhook(request: Request) -> Response:
raw_body = await request.body() # パース前の生のバイト列
try:
event = verify_stripe_event(raw_body, request.headers.get("stripe-signature"), WEBHOOK_SECRET)
except InvalidWebhook:
return Response(status_code=400)
await run_in_threadpool(_handle, event) # 例外は 500 になり、Stripe が再送する
return Response(status_code=200)
DUPLICATE と IGNORED(購読していない種類)も 200 を返します。非 2xx を返すと Stripe は再送を続けるからです。逆に、自分の側で処理できなかったときは必ず非 2xx を返します。「とりあえず 200 を返してからログに残す」実装は、Case A と同じ取りこぼしを生みます。
6. 非同期版:受信の確定と処理を分ける
Stripe の公式ドキュメントは、タイムアウトを起こし得る複雑な処理の前に、素早く 2xx を返すよう推奨しています。業務処理が重い、あるいは外部 API を呼ぶ場合は、受信の確定と処理を分けます。
def receive_async(
conn: psycopg.Connection, *, provider: str, endpoint: str, event: VerifiedEvent
) -> Outcome:
"""受信記録だけをコミットして 2xx を返す。以降の再試行責任はプロバイダからこちらへ移る。"""
with conn.transaction():
inserted = conn.execute(
"""
INSERT INTO webhook_inbox (provider, endpoint, event_id, event_type, payload)
VALUES (%s, %s, %s, %s, %s)
ON CONFLICT (provider, endpoint, event_id) DO NOTHING
RETURNING event_id
""",
(provider, endpoint, event.id, event.type, Jsonb(event.payload)),
).fetchone()
return Outcome.PROCESSED if inserted else Outcome.DUPLICATE
一見 Case A と同じに見えますが、決定的な違いがあります。受信記録が「処理済み」ではなく「未処理(processed_at IS NULL)」として残ることです。Case A では受信記録が処理済みの証拠として扱われたので、処理の失敗が永久消失になりました。非同期版では、受信記録は自分たちが処理する責任を引き受けた証拠です。COMMIT して 2xx を返した時点で、再試行の責任は Stripe からこちらのワーカーに移ります。
ワーカーは「取り出し」と「処理」を別のトランザクションに分けます。
from datetime import timedelta
InboxKey = tuple[str, str, str] # (provider, endpoint, event_id)
CLAIM_INBOX_SQL = """
UPDATE webhook_inbox AS w
SET attempts = w.attempts + 1,
available_at = now() + %(lease)s -- lease: 処理中は他のワーカーから見えなくする
FROM (SELECT provider, endpoint, event_id
FROM webhook_inbox
WHERE processed_at IS NULL
AND available_at <= now()
AND attempts < %(max_attempts)s
ORDER BY available_at
LIMIT 1
FOR UPDATE SKIP LOCKED) AS picked
WHERE (w.provider, w.endpoint, w.event_id) = (picked.provider, picked.endpoint, picked.event_id)
RETURNING w.provider, w.endpoint, w.event_id
"""
def claim_next(
conn: psycopg.Connection, *, max_attempts: int = 10, lease: timedelta = timedelta(minutes=2)
) -> InboxKey | None:
"""未処理を1件選び、試行回数を数えてからコミットする(処理より前に数えるので、
処理中にプロセスごと落ちても回数は残り、上限で止まる)。"""
with conn.transaction():
row = conn.execute(CLAIM_INBOX_SQL, {"lease": lease, "max_attempts": max_attempts}).fetchone()
return (row[0], row[1], row[2]) if row else None
def process_claimed(conn: psycopg.Connection, key: InboxKey, *, timeout: str = "30s") -> bool:
"""取り出した1件を処理する。行をロックして processed_at を確かめ直し、
業務更新と processed_at を同じトランザクションで確定する。
lease が切れて別のワーカーが同じ行を取っても、ロックと再確認があるので業務更新は1回だけ。
(業務の効果が「処理済み」の確認と同じ DB トランザクションにあるので fencing token は要らない。
外部 API を呼ぶ処理では成り立たないので、Outbox Dispatcher の lease + fencing を使う)
"""
try:
with conn.transaction():
# 1件の処理時間に上限を付ける。lease より十分短くする
conn.execute("SELECT set_config('statement_timeout', %s, true)", (timeout,))
row = conn.execute(
"""
SELECT event_type, payload FROM webhook_inbox
WHERE provider = %s AND endpoint = %s AND event_id = %s
AND processed_at IS NULL
FOR UPDATE
""",
key,
).fetchone()
# SKIP LOCKED にしない: 他ワーカーの claim 文が候補として一瞬ロックしただけでも
# 飛ばしてしまい、lease が切れるまで処理されなくなる。待つ時間は statement_timeout が上限
if row is None:
return False # 待っている間に別のワーカーが処理済みにした
event_type, payload = row
apply_event(conn, VerifiedEvent(id=key[2], type=event_type, payload=payload))
conn.execute(
"""
UPDATE webhook_inbox SET processed_at = now(), last_error = NULL
WHERE provider = %s AND endpoint = %s AND event_id = %s
""",
key,
)
return True
except Exception as exc: # ワーカー境界: 1件の失敗でループ全体を止めない
# 業務更新はロールバック済み。失敗を記録し、lease の残りを待たずにバックオフ後へ回す
# (試行回数は claim_next で数え済み)
with conn.transaction():
conn.execute(
"""
UPDATE webhook_inbox
SET last_error = %s,
available_at = now() + make_interval(secs => least(300, 2 ^ attempts) * random())
WHERE provider = %s AND endpoint = %s AND event_id = %s
AND processed_at IS NULL -- 他のワーカーが処理済みにした行は触らない
""",
(type(exc).__name__, *key),
)
return False
def process_next(conn: psycopg.Connection, *, max_attempts: int = 10) -> bool:
"""1件取り出して処理する。取り出せる行が無ければ False(ワーカーは少し待ってから再試行)。"""
key = claim_next(conn, max_attempts=max_attempts)
if key is None:
return False
process_claimed(conn, key)
return True
分ける理由は試行回数を確実に数えるためです。取り出しと処理を1つのトランザクションにすると、処理中にプロセスそのものが落ちたとき(OOM で kill される等)、回数の加算も一緒に ROLLBACK されます。毎回ワーカーを落とすイベントは、上限に達しないまま永遠に再取得されます。claim_next は短いトランザクションで回数を +1 し、available_at を lease の分だけ先へ進めてからコミットします。処理の途中で落ちても回数は残り、lease が切れれば再取得され、上限に達すれば取り出されなくなります(検証済み)。
処理側の process_claimed は、行を FOR UPDATE でロックして processed_at を確かめ直してから業務更新します。lease が切れて別のワーカーが同じ行を取っても、後から来た側はロックを待ち、処理済みを見て何もしません。検証では、取り出した後に止まったワーカー A の lease を切らせ、ワーカー B が処理した後に A を再開させても、元帳は1件でした。
ここで SKIP LOCKED を使わないのには理由があります。別のワーカーの取り出し文が、この行を候補として一瞬ロックすることがあります。SKIP LOCKED だとその瞬間に処理を諦めてしまい、lease が切れるまで行が処理されません。筆者は最初 SKIP LOCKED で書き、4ワーカーの並行テストで実際にこの取りこぼしを観測しました。待つ時間は statement_timeout(lease より十分短くする)が上限なので、待ち続けることはありません。処理が statement_timeout を超えると QueryCanceled として失敗が記録されます(検証済み)。
Dispatcher と違って fencing token が要らないのは、業務の効果(元帳への計上)が「処理済みか」の確認と同じ DB トランザクションの中にあるからです。ロックと再確認だけで、業務更新が2回コミットされることはありません。
ハンドラが外部 API(メール送信、Stripe への書き戻しなど)を呼ぶなら、この前提が崩れます。外部呼び出しは DB のトランザクションに含まれないので、ロックと再確認では二重の呼び出しを防げません。その場合は、外部呼び出しを Outbox に積み、別のディスパッチャで実行します。複数のディスパッチャで同じ行を二重に実行しないための lease と fencing token は、Outbox Dispatcherを複数workerで安全に動かす方法で扱います。
attempts が上限に達した行は、ワーカーに取り出されなくなります。これは人の確認が必要な行です。件数を監視し(9章)、原因を直したら attempts を戻して再処理させます。
7. DynamoDB による重複排除との比較
サーバーレス構成では、DynamoDB の条件付き書き込み(attribute_not_exists)で重複排除するのが定番です。筆者も、DynamoDB の条件付き書き込みによる Webhook の重複排除と、重複排除マーカーを業務更新と同じトランザクションに含める設計を、それぞれ別の本番案件で運用してきました。どちらが正しいかではなく、業務データがどこにあるかで選びます。
| 観点 | DynamoDB 重複排除ストア(業務データは別DB) | DynamoDB(業務データも DynamoDB) | PostgreSQL Transactional Inbox |
|---|---|---|---|
| 原子性 | 無い。マーカーと業務更新は別ストアへの2つの書き込み。Case A か Case B のどちらかが残る | TransactWriteItems でマーカーと業務更新を原子的に書ける(最大100アクション) | 1トランザクションで原子的 |
| 同時重複 | 条件付き書き込みで1本だけ通る | 同左 | 一意インデックスの待ちで直列化。1本目が失敗すれば2本目が引き継ぐ |
| 保持期間 | TTL で自動削除。ただし公式には「通常、期限後数日以内」に削除される(即時ではない) | 同左 | 自分で削除ジョブを書く |
| 失敗の分離 | 重複排除ストアの障害で受信全体が止まる(依存先が増える) | 1つのストア | 業務DBと同じ障害ドメイン(DB が落ちれば業務も止まるので、追加の依存は無い) |
| 負荷 | 業務DBに負荷を掛けない | — | 受信1件ごとに業務DBへ書き込みが増える |
| スループット | 高い。キー単位で水平スケール | 高い | 業務DBの書き込み能力に依存 |
| 運用 | IAM、テーブル、TTL の管理が増える | — | 既存のDB運用に乗る |
筆者の判断基準は次の通りです。
- 業務データが PostgreSQL にある:Transactional Inbox を使う。DynamoDB に重複排除だけを置くと、原子性を失う代わりに得るものが少ない。
- 業務データも DynamoDB にある:
TransactWriteItemsで、マーカーの条件付き Put と業務更新を1つのトランザクションに入れる。原子性は PostgreSQL と同等に得られる。 - 重複排除ストアと業務データが別にならざるを得ない:どちらの失敗を許容するか決める。取りこぼしを避けたいなら「業務処理 → マーカー」の順にし(Case B 側)、二重処理は業務キーの一意制約やハンドラの冪等性で吸収する。
TTL を使う場合は、保持期間をプロバイダの再送期間より長くします。Stripe は最大3日間再送するので、TTL を24時間にすると、3日目に届いた再送は重複として弾かれずに処理されます。TTL は削除の下限を決めるだけで、期限切れのアイテムは数日残ることがあります。期限切れのアイテムを読み取りで除外するなら、その除外条件も保持期間の一部として設計します。
8. プロバイダへの書き戻し:冪等キーとその寿命
ハンドラの中で Stripe の API を呼び返すことがあります(例:invoice.paid を受けて、顧客のメタデータを更新する)。この書き戻しにも冪等性が必要です。
Stripe の API は Idempotency-Key ヘッダーで冪等なリトライをサポートします。公式ドキュメントによれば、同じキーで送られた2回目以降のリクエストには、成功・失敗を問わず最初のリクエストの結果(500 エラーを含む)が返されます。キーは最大255文字で、メールアドレスのような個人情報を含めないよう案内されています。
キーは、イベントIDと操作の種類から決定的に作ります。
idempotency_key = f"webhook:{event.id}:update-customer-metadata"
リクエストごとにランダムなキーを作ると、再送のたびに別のリクエストとして扱われ、冪等性になりません。
ここに見落とされがちな落とし穴があります。Stripe の公式ドキュメントは、キーは少なくとも24時間経過した後にシステムから削除され得て、削除後に同じキーを使うと新しいリクエストとして扱われると説明しています。一方、Webhook は最大3日間再送されます。つまり、3日目に届いた再送で同じキーを使って呼び返すと、Stripe 側では「初めてのリクエスト」として実行される可能性があります。
対策は、書き戻しの結果(作成されたオブジェクトのIDなど)を自分の DB に記録し、記録があれば呼ばないことです。冪等キーは「数時間以内のリトライ」を守る道具であって、「何日後でも1回」を保証する道具ではありません。外部 API の副作用を長期間にわたって1回に収束させる設計は、Exactly-onceは何を保証するのかで詳しく扱います。
9. 運用:監視と保持
-- 未処理の受信と、最も古い未処理の待ち時間
SELECT count(*) AS unprocessed,
max(now() - received_at) AS oldest_age,
count(*) FILTER (WHERE attempts >= 10) AS needs_review
FROM webhook_inbox
WHERE processed_at IS NULL;
-- 直近1時間の受信と重複率(重複は DUPLICATE を返した数をアプリ側で数える)
SELECT event_type, count(*)
FROM webhook_inbox
WHERE received_at > now() - interval '1 hour'
GROUP BY event_type
ORDER BY 2 DESC;
監視する指標:
- 最古の未処理の待ち時間:非同期版で処理が止まっている最初の兆候。
needs_review(上限まで失敗した行)の件数:0件以外ならアラート。- 5xx を返した回数:Stripe のダッシュボードの配信失敗と突き合わせる。
DUPLICATEの割合:急増は、レスポンスが遅すぎて Stripe がタイムアウト扱いしている兆候。
保持期間は、プロバイダの再送期間(Stripe なら3日)より十分長くします。筆者は、再送期間に障害調査の期間を足して、処理済みの行を30日程度で削除することが多いです。
DELETE FROM webhook_inbox
WHERE processed_at IS NOT NULL
AND processed_at < now() - interval '30 days';
それでも Webhook だけに頼ると、Stripe 側が送らなかったイベントや、こちらが長期間落ちていた間のイベントは拾えません。重要な状態は、定期的に Stripe の API と突き合わせて修復します。その設計は Reconciliation(整合性チェック)の設計 で扱います。
10. テスト
PostgreSQL のロックと一意インデックスの挙動を確かめるには、本物の PostgreSQL に複数の接続を張ったテストが必要です。筆者が PostgreSQL 18 のコンテナで実行して通過を確認したケースは次の通りです。
| テスト | 期待する結果 |
|---|---|
| 同じイベントを2回 | 1回目 PROCESSED、2回目 DUPLICATE、元帳1件 |
| 同じイベントを8接続から同時に | PROCESSED 1本、DUPLICATE 7本、元帳1件 |
| 1本目が処理中に失敗し、2本目が同時に待っている | 2本目が1本目の ROLLBACK 後に PROCESSED、元帳1件 |
| 業務処理が失敗 → 再送 | 受信記録は残らず、再送で PROCESSED |
invoice.paid が invoice.created より先に届く | 最終状態は paid |
| 別の Event ID で同じ請求の入金が2回 | 元帳1件 |
void 済みの請求に invoice.paid が届く | 計上されず、通知も積まれない |
| 同じイベントを別エンドポイントで受信 | 両方 PROCESSED |
| Case A(先に記録) | 再送が DUPLICATE になり、元帳0件(失敗の再現) |
| Case B(後で記録、間でクラッシュ) | 元帳2件(失敗の再現) |
| 非同期版:ワーカーの処理失敗 | attempts=1、last_error が残り、再試行で処理される |
| 非同期版:50件を4ワーカーで | 元帳ちょうど50件、未処理0件、どの行も試行1回 |
| 非同期版:取り出し直後にクラッシュを3回(上限3) | attempts=3 で取り出されなくなる(人の確認へ) |
| 非同期版:lease 切れで別ワーカーが処理した後、元のワーカーが戻る | 元のワーカーは何もしない、元帳1件 |
| 非同期版:処理中の行を別ワーカーが処理しようとする | ロックを待ち、処理済みを見て引く |
非同期版:処理が statement_timeout を超える | QueryCanceled として記録され、未処理のまま |
同時実行のテストは threading.Barrier でスレッドの開始をそろえ、各スレッドに専用の接続を持たせます。
def test_concurrent_same_event(dsn, conn):
e = ev("evt_1", "invoice.paid")
barrier = threading.Barrier(8)
def worker():
with psycopg.connect(dsn) as c:
barrier.wait()
return receive_sync(c, provider="stripe", endpoint="billing", event=e)
with ThreadPoolExecutor(8) as ex:
results = [f.result() for f in [ex.submit(worker) for _ in range(8)]]
assert results.count(Outcome.PROCESSED) == 1
assert results.count(Outcome.DUPLICATE) == 7
assert conn.execute("SELECT count(*) FROM ledger_entries").fetchone()[0] == 1
まとめ
- Webhook は同じイベントが複数回、同時に、順不同で届く前提で作る(Stripe は公式に明記)。
- 受信記録を先に別コミットすると永久消失(at-most-once 側)、後で別コミットすると二重処理(at-least-once 側)。原因はどちらも「2つの書き込みを別々にコミットしている」こと。
- Transactional Inbox は、受信記録と業務更新を1トランザクションに入れる。同時重複は一意インデックスの待ちで直列化され、失敗すれば両方消えて再送で再処理される。
- イベントIDの重複排除だけでは足りない。状態遷移として書き、業務キーに一意制約を張り、
createdで順序を判定しない。 - 業務データが PostgreSQL なら Inbox、DynamoDB なら
TransactWriteItems。ストアが分かれるなら、どちらの失敗を許容するかを明示的に選ぶ。 - 冪等キーには寿命がある。数日後の再送でも1回にしたい書き戻しは、結果を自分の DB に記録する。
Stripe の決済実装全体(Checkout、サブスクリプションの状態機械、テスト)は StripeのWebhookと冪等性を本番品質で実装する に、SQS と Lambda で at-least-once の配信を冪等に消費する方法は SQS + Lambda + EventBridge で冪等な非同期処理を作る にまとめています。