メインコンテンツへスキップ
信頼性・非同期・リアルタイム
決済
信頼性
PostgreSQL
Python
アーキテクチャ設計

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

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

公開日
読了時間
28分
著者
友田 陽大
シェア

結論から書きます。 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 を受け取れなければ失敗とみなします。ハンドラは正しく処理したのに、レスポンスがタイムアウトやネットワーク断で届かなければ、同じイベントがもう一度届きます。つまり受信側は、同じイベントが複数回、しかも同時に届く可能性を前提に作る必要があります。


二重課金・Webhookの取りこぼし・支払い失敗による解約でお困りですか?

決済信頼性の診断と、課金まわりの技術選定の技術顧問

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

  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 受信処理本体

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 で冪等な非同期処理を作る にまとめています。

よくある質問

Webhookの冪等性はイベントIDにUNIQUE制約を張れば十分ですか?
不十分です。UNIQUE制約で受信の重複は弾けますが、その記録を業務処理と別のトランザクションでコミットすると、業務処理が失敗したときに再送が重複扱いされてイベントが永久に失われます。受信記録と業務更新を同じトランザクションでコミットし、さらにハンドラ自体を業務キーで冪等にする必要があります。
Transactional Inboxとは何ですか?
受信したメッセージのIDを記録する行(inbox)の挿入と、そのメッセージによる業務更新を、同じデータベーストランザクションでコミットする受信側のパターンです。送信側で業務更新とイベント記録を同じトランザクションに入れるTransactional Outboxと対になり、どちらも副作用と台帳を1つのコミットで原子化します。
Webhookの処理は同期と非同期のどちらで行うべきですか?
業務更新がDB内で短時間に終わるなら、受信記録と同じトランザクションで同期処理してから2xxを返すのが最も単純です。重い処理や外部API呼び出しが要るなら、受信記録だけをコミットして2xxを返し、ワーカーが受信記録の処理済みと業務更新を同じトランザクションで確定させます。どちらでも、受信記録のコミットが再試行責任の境界になります。
StripeのWebhookでevent.createdを使って順序を判定してよいですか?
Stripeの公式ドキュメントは、createdは秒単位で別イベントが同じ値を持ち得るため、順序や処理済みの判定に使わないよう明記しています(2026年10月時点)。順序に依存しない状態遷移として書くか、必要ならStripe APIから最新のオブジェクトを取得して判断します。
DynamoDBで重複排除するのは間違いですか?
間違いではありません。業務データもDynamoDBにあるならTransactWriteItemsで重複排除マーカーと業務更新を原子的に書けます。問題になるのは、重複排除マーカーがDynamoDBにあり業務データがPostgreSQLにあるように、2つのストアをまたぐ場合です。その場合は2つの書き込みを原子的にできないため、どちらかの失敗パターンが残ります。

参考文献

友田

友田 陽大

経済産業大臣賞 受賞プロダクト開発者。TypeScript + Python + AWS で、SaaS・業界DX・実用レベルの生成AI(RAG)を、要件定義からインフラ・運用まで一人で完遂します。

二重課金・Webhookの取りこぼし・支払い失敗による解約でお困りですか?

決済信頼性の診断と、課金まわりの技術選定の技術顧問

「まれに二重課金が出る」「Webhook の取りこぼしに後から気づく」「支払い失敗からの解約が止まらない」——原因は個別のバグより、冪等性・順序保証・リトライ方針・状態機械のどれかが構造的に欠けていることがほとんどです。既存実装の読み取りから、再送で壊れない受信設計、失敗した支払いの回収フロー、自前実装とサービス採用の線引きまで、技術顧問として一緒に決めます。

プロジェクト単位(請負)・技術顧問のどちらにも対応可能です。まずは30分の無料技術相談から。

最短ルート:カレンダーから直接予約

相談内容が固まっている方は、フォーム送信よりその場で日程を確定する方がスムーズです。下記から空き時間をお選びください。

  • 30分のオンライン無料相談
  • Google Meet / Zoom / Microsoft Teams
  • NDA 商談前締結可・無理な営業はいたしません
無料相談の空き枠を予約する

あわせて読みたい