# Outbox Dispatcherを複数workerで安全に動かす：FOR UPDATE SKIP LOCKEDだけでは足りない理由とLease・Fencing Token

> PostgreSQLのOutboxテーブルを複数workerで処理するとき、SELECT ... FOR UPDATE SKIP LOCKEDしてCOMMITした後に外部APIを呼ぶと同じメッセージが二重に処理される理由と、lease・fencing token（lease_version）による正しいclaim/finalizeを解説。バックオフ、dead-letter、集約内の順序、監視SQLまで、PostgreSQL 18で並行実行テストしたPythonコードで示します。

- 公開日: 2026-10-04
- 著者: 友田 陽大
- タグ: PostgreSQL, 信頼性, アーキテクチャ設計, Python, データベース
- URL: https://tomodahinata.com/blog/outbox-dispatcher-skip-locked-lease-fencing-token-guide
- カテゴリ: 信頼性・非同期・リアルタイム
- 総合ガイド: https://tomodahinata.com/blog/transactional-outbox-pattern-reliable-event-publishing-guide

## 要点

- FOR UPDATE の行ロックはCOMMITで解放される。『SELECT ... FOR UPDATE SKIP LOCKED → COMMIT → 外部API』と書くと、外部API実行中に別workerが同じ行を取り、二重に処理する
- SKIP LOCKEDの役割は『同じ瞬間に同じ行を2人が掴まない』ことだけ。COMMIT後も続く所有権は、status・lease_until・lease_versionとしてテーブルに永続化する
- claimで lease_version を+1し、finalizeは『WHERE id = ? AND lease_version = ?』で行う。leaseが失効して他workerが引き継いだら、古いworkerの確定は0行更新で弾かれる（fencing token）
- fencingが守るのはDBの状態だけ。古いworkerがすでに外部APIを呼んでいたら、その副作用は取り消せない。下流の冪等性（メッセージIDによる重複排除）は必須のまま残る
- 失敗は分類して扱う：一時的失敗はジッター付きバックオフ、業務上の保留は試行回数を消費せず延期、恒久的失敗と試行回数超過はdeadにして人が確認する

---

**結論から書きます。** Outbox を複数の worker で処理するとき、`SELECT ... FOR UPDATE SKIP LOCKED` で行を取り、**COMMIT してから外部 API を呼ぶと、同じメッセージを2つの worker が処理します。** 行ロックは COMMIT で解放されるので、外部 API を呼んでいる間、その行はもう誰のものでもないからです。正しい設計では、短いトランザクションで行を `processing` に変え、**所有者・期限（lease）・世代番号（fencing token）をテーブルに書いてから COMMIT** します。確定するときは世代番号が一致する場合だけ書き込み、期限切れで引き継がれた古い worker の書き込みを弾きます。それでも外部の副作用の重複はゼロにならないので、**下流の冪等性は必須のまま**です。

この記事は、[Transactional Outbox](/blog/transactional-outbox-pattern-reliable-event-publishing-guide) で「業務更新とイベントを同じトランザクションで記録する」ところまでできた後の、**配送側（Dispatcher / Relay）を本番で並列化する**ための設計を扱います。コードは **PostgreSQL 18.6 / psycopg 3.3 で、4 worker の並行実行・worker のクラッシュ・lease の失効を再現する11のテストを実行して検証済み**です（2026年10月時点）。

> **用語**：本記事では、Outbox の行を取り出して外部（メッセージブローカーや外部 API）へ送るプロセスを **Dispatcher** と呼びます。microservices.io の Polling Publisher、一般に Relay と呼ばれるものと同じ役割です。

---

## 1. 壊れた Dispatcher：COMMIT の後に外部 API を呼ぶ

よく見かける実装です。

```python
# ❌ 壊れた例：COMMIT でロックが消えた後に外部 I/O をしている。コピーしないでください
def dispatch_broken(conn: psycopg.Connection, publisher: Publisher, *, batch: int = 50) -> list[int]:
    with conn.transaction():
        rows = conn.execute(
            """
            SELECT id, aggregate_id, event_type, payload FROM outbox
             WHERE status = 'pending' ORDER BY id LIMIT %s
               FOR UPDATE SKIP LOCKED
            """,
            (batch,),
        ).fetchall()
    # ← ここで COMMIT 済み。行ロックはもう無い。別 worker が同じ行を SELECT できる
    for id_, agg, typ, payload in rows:
        publisher.publish(message_id=str(id_), key=agg, event_type=typ, payload=payload)
        with conn.transaction():
            conn.execute("UPDATE outbox SET status = 'done', done_at = now() WHERE id = %s", (id_,))
    return [r[0] for r in rows]
```

`FOR UPDATE SKIP LOCKED` を使っているので、複数の worker で動かしても安全に見えます。しかし PostgreSQL の公式ドキュメントにある通り、**行レベルのロックはトランザクションの終了時に解放されます**。

```text
時刻  Worker A                                   Worker B
t1    BEGIN; SELECT ... FOR UPDATE SKIP LOCKED → id=1..3 をロック
t2    COMMIT（ロック解放。status は pending のまま）
t3    publish(id=1) …（外部 API 待ち）
t4                                               BEGIN; SELECT ... SKIP LOCKED → id=1..3 が取れる
t5                                               COMMIT; publish(id=1)
t6    UPDATE status='done'                       UPDATE status='done'
      → id=1..3 がそれぞれ2回発行される
```

検証では、Worker A を publish の途中で止めておき、その間に Worker B を走らせると、B は A と同じ行をすべて取得し、各メッセージが2回発行されました。

### 「トランザクションを開けたまま publish する」も解決にならない

逆に、COMMIT を publish の後にすれば、行ロックは publish の間も続きます（[既存の Outbox 記事](/blog/transactional-outbox-pattern-reliable-event-publishing-guide)のポーリングリレーはこの形です）。行の二重取得は防げますが、別の問題が出ます。

- **外部 API の待ち時間だけ、トランザクションと DB 接続を握り続ける。** ブローカーが遅延すると、接続プールと行ロックが溜まる。
- **publish は成功したのに COMMIT が失敗する**（接続断、フェイルオーバー）と、行は `pending` に戻り、次の worker が再送する。これ自体は at-least-once として許容できるが、長いトランザクションほどこの窓は広がる。
- **長時間のトランザクションは VACUUM を妨げ**、テーブルの肥大化を招く。

小規模で publish が速いうちは、この形でも実用になります。worker を増やす、外部 API が遅い、処理に数秒以上かかる、のどれかに当てはまるなら、次章の lease 方式に切り替えます。

---

## 2. SKIP LOCKED の本当の役割

PostgreSQL の公式ドキュメントは `SKIP LOCKED` を次のように説明しています（要約）。

> `SKIP LOCKED` を指定すると、即座にロックできない行はスキップされる。ロックされた行を飛ばすとデータの一貫しない見え方になるので**汎用的な用途には向かない**が、**キューのようなテーブルに複数の消費者がアクセスする際のロック競合を避ける**のに使える。

つまり `SKIP LOCKED` が保証するのは、**「同じ瞬間に、同じ行を2つのトランザクションが掴まない」**ことだけです。COMMIT の後のことは何も保証しません。

> **row lock ≠ durable ownership**：行ロックはトランザクションの寿命しか持たない。COMMIT 後も「この行は Worker A が処理中だ」という事実を残したいなら、その事実を**テーブルの列として永続化**する必要がある。

これは、`pg_advisory_xact_lock` でも同じです。トランザクションレベルの Advisory Lock も COMMIT で消えるので、COMMIT 後の所有権には使えません（[PostgreSQL Advisory Lock と FOR UPDATE の違い](/blog/postgresql-advisory-lock-pg-advisory-xact-lock-race-condition-guide)）。

---

## 3. 正しい設計：claim → COMMIT → 外部 API → fenced finalize

```text
   pending ──claim（短いTx, lease_version+1）──▶ processing
      ▲                                              │
      │  一時的失敗：backoff して戻す                    │ 外部 API（トランザクションの外）
      │  業務上の保留：試行回数を戻して延期               ▼
      └──────────────── fenced update ◀──────── 結果を分類
                          │      │
                          ▼      ▼
                        done    dead（恒久的失敗・試行回数超過 → 人が確認）

   processing のまま lease_until を過ぎた行 → 次の claim で再取得（lease_version+1）
```

### 3.1 スキーマ

```sql
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(),

    -- ここから下が Dispatcher 用の列
    status         text        NOT NULL DEFAULT 'pending'
                   CHECK (status IN ('pending', 'processing', 'done', 'dead')),
    available_at   timestamptz NOT NULL DEFAULT now(),  -- 次に取り出してよい時刻（バックオフ）
    attempts       integer     NOT NULL DEFAULT 0,      -- claim された回数（クラッシュも1回と数える）
    lease_owner    text,                                -- 観測用。正しさは lease_version が担う
    lease_until    timestamptz,                         -- この時刻を過ぎた processing は再取得可能
    lease_version  bigint      NOT NULL DEFAULT 0,      -- fencing token。claim のたびに +1
    last_error     text,
    done_at        timestamptz
);

CREATE INDEX outbox_ready  ON outbox (available_at, id) WHERE status = 'pending';
CREATE INDEX outbox_leased ON outbox (lease_until)      WHERE status = 'processing';
```

既存の Outbox 記事のスキーマ（`published_at` が NULL なら未発行）に、claim のための列を足した形です。

- **`lease_version` が fencing token です。** claim のたびに +1 され、確定（finalize）のときに一致を確認します。`lease_owner` は「誰が持っているか」を調べるための観測用の列で、正しさは `lease_version` だけで保証されます。owner だけで判定すると、同じ owner 名の worker が再起動したときに古い処理と新しい処理を区別できません。
- **`attempts` は claim の時点で増やします。** 処理を始めたとたんに worker がクラッシュする「毒メッセージ」も試行回数として数え、無限に再取得されるのを防ぐためです。
- **時刻はすべて DB の `now()` で計算します。** worker ごとの時計のずれで lease の判定が狂わないようにするためです。

### 3.2 claim：短いトランザクションで所有権を書き込む

```python
from dataclasses import dataclass
from datetime import timedelta
from typing import Any

import psycopg
from psycopg.rows import class_row


@dataclass(frozen=True)
class Claimed:
    id: int
    lease_version: int
    attempts: int
    aggregate_id: str
    event_type: str
    payload: dict[str, Any]


CLAIM_SQL = """
WITH picked AS (
    SELECT id
      FROM outbox
     WHERE (status = 'pending'    AND available_at <= now())
        OR (status = 'processing' AND lease_until  <  now()     -- 落ちた worker の取り残し
            AND attempts < %(max_attempts)s)
     ORDER BY id
     LIMIT %(batch)s
       FOR UPDATE SKIP LOCKED
)
UPDATE outbox AS o
   SET status        = 'processing',
       lease_owner   = %(owner)s,
       lease_until   = now() + %(lease)s,
       lease_version = o.lease_version + 1,
       attempts      = o.attempts + 1
  FROM picked
 WHERE o.id = picked.id
RETURNING o.id, o.lease_version, o.attempts, o.aggregate_id, o.event_type, o.payload
"""


# claim のたびに worker を落とし続ける毒メッセージは確定まで辿り着かない。試行回数を使い切ったまま lease 切れになった行は dead に送る
REAP_SQL = """
UPDATE outbox
   SET status = 'dead', last_error = 'lease expired after max attempts (worker crash?)',
       lease_owner = NULL, lease_until = NULL
 WHERE status = 'processing' AND lease_until < now() AND attempts >= %(max_attempts)s
"""


def claim(
    conn: psycopg.Connection, *, owner: str, batch: int, lease: timedelta, max_attempts: int = 12
) -> list[Claimed]:
    params = {"batch": batch, "owner": owner, "lease": lease, "max_attempts": max_attempts}
    with conn.transaction(), conn.cursor(row_factory=class_row(Claimed)) as cur:
        cur.execute(REAP_SQL, params)  # 試行回数を使い切ったまま lease 切れの行を dead へ
        cur.execute(CLAIM_SQL, params)
        return sorted(cur.fetchall(), key=lambda c: c.id)
```

この1文が保証すること：

1. **同時に claim した worker 同士は、互いに違う行を取る。** `SKIP LOCKED` がここで効く。検証では、200件を4つの worker が7件ずつ同時に claim し続け、取得した ID は重複なしでちょうど200件だった。
2. **COMMIT 後も「processing・所有者・期限・世代」がテーブルに残る。** 他の worker の claim は `status = 'pending'` か「lease 切れの `processing`」しか取らないので、lease 中の行には触れない。
3. **落ちた worker の行は、lease が切れれば自動的に再取得される。** 再取得のたびに `lease_version` が上がる。
4. **worker を落とし続けるメッセージは、試行回数の上限で止まる。** 確定まで辿り着かない毒メッセージは例外処理の分岐を一度も通らないので、claim の側で「上限に達したまま lease が切れた行」を `dead` に送る（`REAP_SQL`）。検証では、claim して確定しない（＝毎回クラッシュする）ことを3回繰り返すと、`attempts=3` で `dead` になり、それ以上取り出されなかった。

`RETURNING` の行の順序は保証されないので、`id` で並べ直しています。

### 3.3 外部 API はトランザクションの外で呼ぶ

claim のトランザクションは即座に COMMIT されます。外部 API はその後、**トランザクションを開いていない状態**で呼びます。DB 接続もロックも握らないので、外部 API がどれだけ遅くても DB には影響しません。

lease の長さは、**外部 API のタイムアウトより十分長く**します。例えば publish のタイムアウトが10秒なら、lease は60秒にします。lease より長くかかる処理がある場合は、処理の途中で lease を延長する（heartbeat）ことになりますが、Outbox の publish のような短い処理では、タイムアウトを lease より短く設定する方が単純です。

### 3.4 fenced finalize：世代番号が一致する worker だけが確定できる

```python
def _fenced_update(conn: psycopg.Connection, msg: Claimed, set_clause: str, params: dict[str, Any]) -> bool:
    # set_clause はこのモジュール内の定数だけを渡す（外部入力を連結しない）
    with conn.transaction():
        cur = conn.execute(
            f"""
            UPDATE outbox
               SET {set_clause},
                   lease_owner = NULL,
                   lease_until = NULL
             WHERE id = %(id)s
               AND status = 'processing'
               AND lease_version = %(version)s
            """,
            {"id": msg.id, "version": msg.lease_version, **params},
        )
        return cur.rowcount == 1


def mark_done(conn: psycopg.Connection, msg: Claimed) -> bool:
    """False = 自分の lease は失効し、誰かが引き継いだ。"""
    return _fenced_update(conn, msg, "status = 'done', done_at = now(), last_error = NULL", {})
```

`WHERE lease_version = %(version)s` が fencing です。次のタイムラインで効きます。

```text
時刻  Worker A                                  Worker B
t1    claim → id=1, lease_version=1, lease 60秒
t2    publish(id=1) … GC 停止・ネットワーク遅延で 90秒止まる
t3    （lease 失効）
t4                                              claim → id=1, lease_version=2
t5                                              publish(id=1) → 成功
t6                                              mark_done(version=2) → 1行更新
t7    publish 完了 → mark_done(version=1) → 0行更新（弾かれる）
```

検証では、lease を 200 ミリ秒にして claim した Worker A を放置すると、lease 中は Worker B の claim が空で、失効後は B が同じ ID を `lease_version` +1・`attempts` 2 で取得し、A の `mark_done` は `False`、B の `mark_done` は `True` になりました。

**lease が切れていても、誰も引き継いでいなければ、元の worker の確定は成功します**（`lease_until` を WHERE に入れていないため）。期限切れを理由に、成功した処理をわざわざ捨てて再送する必要はないからです。正しさを決めるのは「時刻」ではなく「世代が変わったかどうか」です。

ただし最後の試行だけは例外です。その claim で `attempts` が `max_attempts` に達していて、publish が lease より長引いた場合、誰も引き継いでいなくても、次の `claim()`（どの worker のものでも）が `REAP_SQL` でその行を `dead` にします。元の worker の `mark_done` は行がもう `processing` ではないため `False` を返し、`lost lease` を記録します。メッセージは届いていますが、行は人が確認するために `dead` に残ります。最後の試行でこうならないよう、publish のタイムアウトは lease より十分短くしてください。

### 3.5 fencing が守らないもの

上のタイムラインで、**メッセージ id=1 は2回 publish されています**（t2 の A と t5 の B）。fencing token が防いだのは、A の古い結果で DB の状態を上書きすることだけです。A がすでに外部に送ったメッセージは取り消せません。

Martin Kleppmann が分散ロックの議論で指摘している通り、fencing token で外部の副作用まで守るには、**書き込み先のシステム自身が token を検証して古いものを拒否する**必要があります。メッセージブローカーや多くの外部 API はそうした機能を持たないので、Outbox では代わりに次の2つを組み合わせます。

- **メッセージ ID（`outbox.id`）を、下流の重複排除キーとして渡す。** consumer はこの ID で重複を弾く（[Transactional Inbox でイベントを冪等に受信する方法](/blog/webhook-idempotency-transactional-inbox-postgresql-guide)）。
- **外部 API を直接呼ぶなら、`outbox.id` から作った冪等キーを付ける。** 同じキーの再送は、相手側で1回の効果にまとめられる。

**lease と fencing は「二重処理の頻度を下げ、DB の状態を壊さない」仕組みであって、「二重処理をゼロにする」仕組みではありません。** この区別が、Dispatcher の設計で最も誤解されやすい点です。

---

## 4. 失敗を分類する

外部 API の失敗を一律にリトライすると、直らない失敗でキューが詰まり、すぐ直る失敗を必要以上に待たせます。結果を4つに分類します。

| 分類 | 例 | 処理 | 試行回数 |
| --- | --- | --- | --- |
| 成功・終端の no-op | 送信成功／相手側がすでに目的の状態 | `done` | — |
| 一時的失敗（technical retry） | タイムアウト、5xx、429、接続断 | ジッター付きバックオフで `pending` に戻す。上限を超えたら `dead` | 消費する |
| 業務上の保留（business defer） | 依存するリソース（顧客アカウントなど）がまだ作られていない | 一定時間後に `pending` に戻す | **消費しない** |
| 恒久的失敗 | 4xx の検証エラー、存在しない宛先 | `dead`（人が確認する） | — |

```python
import random


class RetryableError(Exception):
    """一時的な失敗（タイムアウト、5xx、429）。バックオフして再試行する。"""


class PermanentError(Exception):
    """再試行しても直らない失敗（4xx の検証エラー等）。dead にして人が見る。"""


class Defer(Exception):
    """業務上まだ実行できない。試行回数を消費せずに後回しにする。"""

    def __init__(self, delay: timedelta, reason: str) -> None:
        super().__init__(reason)
        self.delay = delay


def reschedule(conn: psycopg.Connection, msg: Claimed, *, delay: timedelta, error: str, refund_attempt: bool) -> bool:
    return _fenced_update(
        conn,
        msg,
        "status = 'pending', available_at = now() + %(delay)s, last_error = %(error)s,"
        " attempts = attempts - %(refund)s",
        {"delay": delay, "error": error[:500], "refund": 1 if refund_attempt else 0},
    )


def mark_dead(conn: psycopg.Connection, msg: Claimed, *, error: str) -> bool:
    return _fenced_update(conn, msg, "status = 'dead', last_error = %(error)s", {"error": error[:500]})


def full_jitter_backoff(attempt: int, *, base: float = 1.0, cap: float = 300.0) -> timedelta:
    """AWS Architecture Blog の "Full Jitter": sleep = random(0, min(cap, base * 2^attempt))"""
    return timedelta(seconds=random.uniform(0, min(cap, base * 2**attempt)))
```

**ジッターを入れる理由**：障害から復旧した瞬間に、溜まっていた全メッセージが同じ間隔でリトライすると、復旧したばかりの外部サービスに負荷が集中します。AWS Architecture Blog の "Full Jitter" は、待ち時間を 0 から上限までの一様乱数にすることで、リトライの時刻を分散させます。リトライ戦略全般は [リトライ・指数バックオフ＋ジッター・サーキットブレーカー実装ガイド](/blog/retry-backoff-circuit-breaker-resilience-patterns-guide) で詳しく扱っています。

**業務上の保留で試行回数を戻す理由**：claim で `attempts` を +1 しているので、保留のたびに回数を消費すると、依存リソースの作成が遅いだけで `dead` になってしまいます。保留は「失敗」ではないので、回数を戻します。ただし、保留が永遠に続く可能性があるなら、`created_at` からの経過時間に上限を設けて `dead` にします。

### 4.1 ディスパッチループ

```python
import logging

log = logging.getLogger("outbox.dispatcher")


def dispatch_once(
    conn: psycopg.Connection,
    publisher: Publisher,
    *,
    owner: str,
    batch: int = 50,
    lease: timedelta = timedelta(seconds=60),
    max_attempts: int = 12,
) -> int:
    """1バッチを処理する。外部呼び出しはトランザクションの外で行う。"""
    messages = claim(conn, owner=owner, batch=batch, lease=lease, max_attempts=max_attempts)
    for msg in messages:
        try:
            publisher.publish(
                message_id=str(msg.id),  # 下流の重複排除キー。再試行しても変わらない
                key=msg.aggregate_id,
                event_type=msg.event_type,
                payload=msg.payload,
            )
        except Defer as d:
            fenced = reschedule(conn, msg, delay=d.delay, error=f"deferred: {d}", refund_attempt=True)
        except RetryableError as exc:
            if msg.attempts >= max_attempts:
                fenced = mark_dead(conn, msg, error=f"retries exhausted: {exc!r}")
            else:
                fenced = reschedule(
                    conn, msg, delay=full_jitter_backoff(msg.attempts), error=repr(exc), refund_attempt=False
                )
        except PermanentError as exc:
            fenced = mark_dead(conn, msg, error=repr(exc))
        else:
            fenced = mark_done(conn, msg)

        if not fenced:
            # publish は済んでいる可能性がある。引き継いだ worker も publish するので重複は下流で吸収する
            log.warning("lost lease", extra={"outbox_id": msg.id, "lease_version": msg.lease_version})
    return len(messages)
```

`Publisher` は、外部ブローカーや API を包む薄い境界です。SDK の例外を `RetryableError` / `PermanentError` / `Defer` に変換するのはこの境界の責任で、ループは SDK の例外を知りません。**分類できない例外はあえて捕まえていません。** 想定外の例外でループが止まれば、プロセスの再起動後に lease の失効で同じ行が再取得され、`attempts` が上限に達すれば `dead` になります。想定外のバグを「一時的失敗」として黙ってリトライし続けるより、表に出る方が安全です。

バッチの末尾の方のメッセージは、先頭のメッセージの publish が遅いと、処理を始める前に lease が切れることがあります。**`batch × 1件あたりの最大処理時間` が lease より十分短くなるように**バッチサイズを決めます。

---

## 5. 集約内の順序が必要な場合

複数の worker で並列に処理すると、同じ注文の `OrderCreated`（id=1）と `OrderPaid`（id=2）を別々の worker が取り、id=2 が先に発行されることがあります。`ORDER BY id` は「取り出す順序」を決めるだけで、「発行が完了する順序」は保証しません。

順序が必要なら、**各集約の先頭のメッセージだけを claim の対象にします。**

> **前提**：この方法は「同じ集約の中では、`id` の順序がコミットの順序と一致する」ことを前提にしています。`id` は INSERT の時点で採番されるので、同じ集約に2つのトランザクションが並行して書くと、id=10 の T1 より id=11 の T2 が先にコミットし、T1 がコミットするまでの間に id=11 が「先頭」として取り出されることがあります。生産者側で、**集約の行を `FOR UPDATE` でロックしてから outbox に INSERT する**（[Transactional Outbox の記事](/blog/transactional-outbox-pattern-reliable-event-publishing-guide)の `with_for_update=True` の例）と、同じ集約への書き込みはコミットまで直列化され、この前提が成り立ちます。

```sql
WITH picked AS (
    SELECT o.id
      FROM outbox AS o
     WHERE ((o.status = 'pending' AND o.available_at <= now())
         OR (o.status = 'processing' AND o.lease_until < now() AND o.attempts < %(max_attempts)s))
       AND NOT EXISTS (                       -- 同じ集約に、未完了の先行メッセージがあれば待つ
            SELECT 1 FROM outbox AS prev
             WHERE prev.aggregate_id = o.aggregate_id
               AND prev.id < o.id
               AND prev.status IN ('pending', 'processing', 'dead'))
     ORDER BY o.id
     LIMIT %(batch)s
       FOR UPDATE SKIP LOCKED
)
UPDATE outbox AS o
   SET status = 'processing', lease_owner = %(owner)s, lease_until = now() + %(lease)s,
       lease_version = o.lease_version + 1, attempts = o.attempts + 1
  FROM picked
 WHERE o.id = picked.id
RETURNING o.id, o.lease_version, o.attempts, o.aggregate_id, o.event_type, o.payload
```

- 同じ集約の id=1 が `pending` か `processing` の間、id=2 は claim の候補にならない。2つの worker が同時に claim しても、id=1 は `SKIP LOCKED` で片方だけが取り、もう片方は id=2 を「先行が未完了」として除外する。
- **`dead` も先行メッセージとして扱う**ので、id=1 が `dead` になると、その集約の後続はすべて止まる。順序を守るなら、壊れたメッセージを飛ばして先に進むことはできないからです。人が id=1 を直すか、意図的に飛ばす（`done` にする）判断をするまで止めます。
- 検証では、同じ集約の30件を4つの worker で処理し、確定の順序が id の昇順と完全に一致しました。
- **保証されるのは claim と確定の順序であり、ブローカーに届く順序ではありません。** lease が失効して引き継がれ、引き継いだ worker が id=1 を確定して id=2 に進んだ後に、古い worker の飛行中の publish（id=1）が届くと、ブローカーには id=2 が id=1 より先に見えることがあります。fencing が守るのは DB の状態だけなので、順序を最後まで安全にするには、受け手側の状態遷移を冪等にして（バージョンや連番のチェックなど）おく必要があります。

代償もあります。集約ごとに直列になるので、1つの集約にメッセージが集中するとスループットが落ちます。`NOT EXISTS` のための `(aggregate_id, id)` のインデックスも必要です。順序が要らないイベント（通知など）まで同じ claim にしないよう、**順序が必要なイベント種別だけ**この claim を使うのが現実的です。下流で順序を吸収できる（状態遷移として冪等に書ける）なら、そもそも順序保証は要りません。

---

## 6. 監視

```sql
-- queue depth（状態別の件数）
SELECT status, count(*) FROM outbox GROUP BY status;

-- 最古の配送待ちメッセージの年齢（最重要。Dispatcher が止まると単調増加する）
SELECT now() - min(created_at) AS oldest_pending_age
  FROM outbox
 WHERE status IN ('pending', 'processing');

-- lease が切れたまま放置されている processing（worker のクラッシュの兆候）
SELECT count(*) AS expired_leases
  FROM outbox
 WHERE status = 'processing' AND lease_until < now();

-- 試行回数の分布（リトライが増えている兆候）
SELECT attempts, count(*)
  FROM outbox
 WHERE status IN ('pending', 'processing')
 GROUP BY attempts ORDER BY attempts;

-- 人の確認が必要なメッセージ
SELECT id, aggregate_id, event_type, attempts, last_error, created_at
  FROM outbox
 WHERE status = 'dead'
 ORDER BY id;
```

アラートの目安：

- **最古の配送待ちの年齢**が SLO（例：5分）を超えた → Dispatcher の停止、外部サービスの障害、スループット不足のどれか。
- **`dead` が1件でも増えた** → 人が原因を確認する。順序付き claim を使っているなら、その集約の後続が止まっている。
- **`expired_leases` が増え続ける** → worker がクラッシュしているか、lease が処理時間に対して短すぎる。
- アプリ側のメトリクスとして **`lost lease` の警告数** → lease の長さかバッチサイズの見直しが必要。

`done` の行は定期的に削除します。削除の保持期間は、障害調査で配送履歴を確認したい期間に合わせます。

```sql
DELETE FROM outbox WHERE status = 'done' AND done_at < now() - interval '7 days';
```

---

## 7. テスト：worker の並行・クラッシュ・失効を再現する

以下は、PostgreSQL 18 のコンテナに対して実行し、すべて通過を確認したテストです。

| テスト | 期待する結果 |
| --- | --- |
| 4 worker が200件を同時に claim | 取得 ID は重複なしで200件 |
| 壊れた Dispatcher（COMMIT 後に publish）を2 worker で | 同じメッセージが2回発行される（失敗の再現） |
| 4 worker で100件を `dispatch_once` | 100件がそれぞれ1回ずつ発行され、全件 `done` |
| claim 後に worker が停止 → lease 失効 | 他 worker が `lease_version` +1 で再取得。古い worker の `mark_done` は `False` |
| lease 失効後だが誰も引き継いでいない | 元の worker の `mark_done` は成功 |
| 一時的失敗を3回（上限3） | `pending`（`attempts=1`）→ … → `dead`（`attempts=3`） |
| 業務上の保留 | `pending` のまま、`attempts=0`、`available_at` が延期される |
| 恒久的失敗 | 即 `dead` |
| claim 後に確定しないことを3回（上限3） | `attempts=3` で `dead`、以後は取り出されない |
| 順序付き claim：集約X（3件）とY（2件） | 1回目は各集約の先頭だけ。先頭が `done` になると次が取れる |
| 順序付き claim：同じ集約30件を4 worker で | 確定の順序が id の昇順と一致 |

クラッシュのテストは、プロセスを殺す必要はありません。**claim した後に確定を呼ばない**ことで、クラッシュと同じ状態（`processing` のまま lease が切れる）を作れます。

```python
def test_crashed_worker_lease_expires_and_stale_finalize_is_fenced(dsn):
    seed(1)
    with psycopg.connect(dsn) as ca, psycopg.connect(dsn) as cb:
        [a] = claim(ca, owner="A", batch=10, lease=timedelta(milliseconds=200))
        assert claim(cb, owner="B", batch=10, lease=timedelta(seconds=30)) == []  # lease 中は取れない
        time.sleep(0.3)
        [b] = claim(cb, owner="B", batch=10, lease=timedelta(seconds=30))
        assert b.id == a.id and b.lease_version == a.lease_version + 1
        assert mark_done(ca, a) is False  # 失効した A の確定は弾かれる
        assert mark_done(cb, b) is True
```

---

## 8. 使わない方がよい場面・別の選択肢

- **配送量が非常に多い、またはメッセージの保持・再生（リプレイ）が要件**：Outbox は DB からブローカーへの「確実な受け渡し」に徹し、配送そのものは Kafka や SQS に任せる。CDC（Debezium など）で WAL から読む方式も、ポーリングの負荷を避けられる（[Transactional Outbox の Polling と CDC の比較](/blog/transactional-outbox-pattern-reliable-event-publishing-guide)）。
- **worker が1つで十分**：単一の worker なら、lease も fencing も要らない。ただし「1つしか動いていない」ことを保証する仕組み（デプロイ時の重複起動を含む）が別に必要になるので、最初から lease 方式にしておく方が安全なことが多い。
- **長時間のジョブ（数分〜数時間）**：lease の heartbeat、キャンセル、進捗管理が必要になり、専用のジョブ基盤（ワークフローエンジンなど）の方が適していることが多い。
- **ポーリングの遅延が許容できない**：`LISTEN/NOTIFY` で新しい行の挿入を worker に知らせ、ポーリングは取りこぼし対策として残す方法がある。ただし PgBouncer の transaction モードでは `LISTEN` が使えないので、専用の接続が必要。

---

## まとめ

- 行ロックは COMMIT で消える。「SKIP LOCKED で取って COMMIT → 外部 API」は二重処理になる。
- SKIP LOCKED の役割は「同じ瞬間に同じ行を2人が掴まない」ことだけ。COMMIT 後の所有権は status・lease_until・lease_version としてテーブルに書く。
- claim で lease_version を +1 し、確定は `WHERE id = ? AND lease_version = ?`。失効して引き継がれた古い worker の確定は0行更新で弾かれる。
- fencing が守るのは DB の状態だけ。外部の副作用の重複はゼロにならないので、メッセージ ID による下流の重複排除と、外部 API の冪等キーを併用する。
- 失敗は4分類（成功／一時的失敗／業務上の保留／恒久的失敗）で扱い、ジッター付きバックオフと dead-letter を用意する。
- 監視の最重要指標は「最古の配送待ちメッセージの年齢」。

Dispatcher が at-least-once で配送したメッセージを、受け取った側で1回の効果にまとめる考え方は [Exactly-onceは何を保証するのか](/blog/exactly-once-at-least-once-idempotency-effectively-once-guide) で、配送が長期間止まったときに DB と外部の状態のズレを検出して直す仕組みは [Reconciliation（整合性チェック）の設計](/blog/reconciliation-architecture-drift-detection-repair-guide) で扱います。
