# Webhookの冪等性をTransactional Inboxで実装する：イベントIDの重複排除だけでは防げない『永久消失』と『二重処理』

> Webhookの冪等性を、受信記録と業務更新を同じPostgreSQLトランザクションでコミットするTransactional Inboxで実装する方法。先に処理済みを記録するとイベントが永久に失われ、後で記録すると二重処理になる理由、同時重複配信・順不同・別Eventオブジェクトへの対処、DynamoDB重複排除との比較を、PostgreSQL 18で検証したPythonコードで解説します。

- 公開日: 2026-10-04
- 著者: 友田 陽大
- タグ: 決済, 信頼性, PostgreSQL, Python, アーキテクチャ設計
- URL: https://tomodahinata.com/blog/webhook-idempotency-transactional-inbox-postgresql-guide
- カテゴリ: 信頼性・非同期・リアルタイム
- 総合ガイド: https://tomodahinata.com/blog/transactional-outbox-pattern-reliable-event-publishing-guide

## 要点

- Webhookは同じイベントが複数回届く前提で設計する。Stripeは本番で最大3日間再送し、同じイベントが複数回届くことや、配信順が生成順と一致しないことを公式に明記している
- 『処理済み』を業務処理の前に別トランザクションでコミットすると、業務処理が失敗したときに再送が重複扱いされ、イベントが永久に失われる（at-most-once側の失敗）
- 業務処理をコミットしてから処理済みを記録すると、その間にプロセスが落ちたとき再送で二重処理になる（at-least-once側の失敗）
- 解は、受信記録（INSERT ... ON CONFLICT DO NOTHING）と業務更新を同じトランザクションに入れるTransactional Inbox。失敗すれば両方消えるので再送で正しく再処理され、同時重複は一意インデックスの待ちで直列化される
- イベントIDの重複排除だけでは足りない。別のEventオブジェクトとして届く同じ業務事実や順不同には、ハンドラ自体を状態遷移と業務キーの一意制約で冪等にして対処する

---

**結論から書きます。** 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 の仕様は[公式ドキュメント](https://docs.stripe.com/webhooks)の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：先に「処理済み」をコミット → イベントの永久消失

```python
# ❌ 壊れた例：受信記録を先に別トランザクションでコミットしている
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
```

```text
1回目の配信:  受信記録 COMMIT → 業務処理で DB エラー → ROLLBACK → 500 を返す
2回目の配信:  受信記録の INSERT が衝突 → 「重複」と判断 → 200 を返す
→ 入金は一度も計上されず、Stripe はもう再送しない
```

これは **at-most-once（高々1回）側の失敗**です。重複は起きない代わりに、取りこぼしが起きます。検証でも、1回目の業務処理を失敗させると、2回目は `DUPLICATE` を返し、元帳は0件のままでした。決済では、二重課金と同じくらい深刻な障害です。

### Case B：業務処理をコミットしてから記録 → 二重処理

```python
# ❌ 壊れた例：業務処理のコミットと受信記録が別トランザクション
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
```

```text
1回目の配信:  業務処理 COMMIT（入金を計上）→ 受信記録の前にクラッシュ → 2xx は返らない
2回目の配信:  受信記録が無い → 未処理と判断 → 業務処理をもう一度 COMMIT
→ 入金が2回計上される
```

これは **at-least-once（少なくとも1回）側の失敗**です。検証でも、受信記録の直前で処理を中断させてから再送すると、元帳は2件になりました。`SELECT` してから `INSERT` しているので、同時に2本届いた場合も両方が「未処理」と判断して二重処理になります（TOCTOU）。

**2つのパターンの原因は同じです。** 受信記録と業務更新という**2つの書き込みを、別々にコミットしている**ことです。どちらを先にしても、間で落ちれば片方だけが残ります。

---

## 3. 解：受信記録と業務更新を1つのトランザクションに入れる（Transactional Inbox）

受信記録の `INSERT` と業務更新を、**同じトランザクション**でコミットします。

```text
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](/blog/transactional-outbox-pattern-reliable-event-publishing-guide) と鏡像の関係にあります。

| | Transactional Outbox（送信側） | Transactional Inbox（受信側） |
| --- | --- | --- |
| 同じトランザクションに入れるもの | 業務更新 ＋ 送るべきイベント | 受け取ったイベントID ＋ 業務更新 |
| 防ぐ失敗 | 業務は更新したのにイベントが出ない／出たのに業務は無い | 受け取ったのに処理されない／2回処理される |
| 残る性質 | 発行は at-least-once（重複発行はあり得る） | 重複配信を吸収し、業務効果を1回に収束させる |

なお、Inbox は**データベースが1つの場合にだけ**この原子性を持ちます。受信記録と業務データが別のストアにあると、また2つの書き込みに戻ります（7章の DynamoDB との比較で扱います）。

---

## 4. スキーマ

```sql
-- 受信台帳。主キーが「同じ配信は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 署名検証：生のボディで行う

```python
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 ハンドラ自体を冪等・順序非依存に書く

```python
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つです。

1. **状態遷移として書く**：`invoice.created` は「無ければ作る」だけで、既存の状態を変えません。`invoice.paid` は `paid` / `void` 以外のときだけ `paid` に進めます。だから `invoice.paid` が `invoice.created` より先に届いても、最終状態は `paid` になります（検証済み）。
2. **`created` で順序を判定しない**：Stripe の公式ドキュメントは、`created` が秒単位で別のイベントが同じ値を持ち得るため、**順序や処理済みの判定に `created` を使わないよう明記**しています（2026年10月時点）。順序が本当に必要な判断（サブスクリプションの最新状態など）は、Stripe API から最新のオブジェクトを取得して決めます。公式ドキュメントも、先に届いたイベントの情報から関連オブジェクトを API で取得する方法を案内しています。
3. **業務キーの一意制約で効果を1回にする**：`ledger_entries` の `UNIQUE (invoice_id, kind)` があるので、別の Event オブジェクト（`evt_A` と `evt_B`）で同じ請求の入金が届いても、計上は1回です（検証済み）。
4. **請求書が `paid` になったときだけ計上する**：`void` 済みの請求書に `invoice.paid` が届いた場合、upsert は状態を変えません。そのまま元帳に計上すると「無効な請求への入金」が記録されるので、状態を読み直して `paid` でなければ計上しません。この組み合わせは異常系なので、Reconciliation で人の確認に回します（検証済み）。

### 5.3 受信処理本体

```python
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 の核心です。

```text
時刻  配信 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 の例）

```python
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 を呼ぶ場合は、受信の確定と処理を分けます。

```python
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 からこちらのワーカーに移ります。

ワーカーは「取り出し」と「処理」を別のトランザクションに分けます。

```python
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で安全に動かす方法](/blog/outbox-dispatcher-skip-locked-lease-fencing-token-guide)で扱います。

`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と操作の種類から決定的に作ります**。

```python
idempotency_key = f"webhook:{event.id}:update-customer-metadata"
```

リクエストごとにランダムなキーを作ると、再送のたびに別のリクエストとして扱われ、冪等性になりません。

ここに見落とされがちな落とし穴があります。Stripe の公式ドキュメントは、**キーは少なくとも24時間経過した後にシステムから削除され得て、削除後に同じキーを使うと新しいリクエストとして扱われる**と説明しています。一方、Webhook は最大3日間再送されます。つまり、**3日目に届いた再送で同じキーを使って呼び返すと、Stripe 側では「初めてのリクエスト」として実行される**可能性があります。

対策は、書き戻しの結果（作成されたオブジェクトのIDなど）を**自分の DB に記録し、記録があれば呼ばない**ことです。冪等キーは「数時間以内のリトライ」を守る道具であって、「何日後でも1回」を保証する道具ではありません。外部 API の副作用を長期間にわたって1回に収束させる設計は、[Exactly-onceは何を保証するのか](/blog/exactly-once-at-least-once-idempotency-effectively-once-guide)で詳しく扱います。

---

## 9. 運用：監視と保持

```sql
-- 未処理の受信と、最も古い未処理の待ち時間
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日程度で削除することが多いです。

```sql
DELETE FROM webhook_inbox
 WHERE processed_at IS NOT NULL
   AND processed_at < now() - interval '30 days';
```

それでも Webhook だけに頼ると、Stripe 側が送らなかったイベントや、こちらが長期間落ちていた間のイベントは拾えません。重要な状態は、定期的に Stripe の API と突き合わせて修復します。その設計は [Reconciliation（整合性チェック）の設計](/blog/reconciliation-architecture-drift-detection-repair-guide) で扱います。

---

## 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` でスレッドの開始をそろえ、各スレッドに**専用の接続**を持たせます。

```python
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と冪等性を本番品質で実装する](/blog/stripe-payments-production-guide-webhooks-idempotency-subscriptions) に、SQS と Lambda で at-least-once の配信を冪等に消費する方法は [SQS + Lambda + EventBridge で冪等な非同期処理を作る](/blog/aws-sqs-lambda-eventbridge-idempotent-async-processing-guide) にまとめています。
