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

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コードで示します。

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

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

この記事は、Transactional Outbox で「業務更新とイベントを同じトランザクションで記録する」ところまでできた後の、配送側(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 を呼ぶ

よく見かける実装です。

# ❌ 壊れた例: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 の公式ドキュメントにある通り、行レベルのロックはトランザクションの終了時に解放されます。

時刻  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 記事のポーリングリレーはこの形です)。行の二重取得は防げますが、別の問題が出ます。

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

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


この記事の実装を、案件として承ります

ORM選定・データモデル設計・ゼロダウンタイム移行を、設計から実装まで承ります

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 の違い)。


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

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

   processing のまま lease_until を過ぎた行 → 次の claim で再取得(lease_version+1)

3.1 スキーマ

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:短いトランザクションで所有権を書き込む

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 だけが確定できる

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 です。次のタイムラインで効きます。

時刻  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 でイベントを冪等に受信する方法)。
  • 外部 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(人が確認する)—
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 から上限までの一様乱数にすることで、リトライの時刻を分散させます。リトライ戦略全般は リトライ・指数バックオフ+ジッター・サーキットブレーカー実装ガイド で詳しく扱っています。

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

4.1 ディスパッチループ

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 の記事の with_for_update=True の例)と、同じ集約への書き込みはコミットまで直列化され、この前提が成り立ちます。

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. 監視

-- 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 の行は定期的に削除します。削除の保持期間は、障害調査で配送履歴を確認したい期間に合わせます。

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_once100件がそれぞれ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 が切れる)を作れます。

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 の比較)。
  • 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は何を保証するのか で、配送が長期間止まったときに DB と外部の状態のズレを検出して直す仕組みは Reconciliation(整合性チェック)の設計 で扱います。

よくある質問

FOR UPDATE SKIP LOCKEDを使えばOutboxを複数workerで処理しても二重処理になりませんか?
トランザクションを外部APIの完了まで開けておけば行ロックは続きますが、接続を握り続け、外部APIの成功後にCOMMITが失敗する問題が残ります。外部APIの前にCOMMITすると行ロックは消え、別workerが同じ行を取れます。どちらの形でも、COMMIT後の所有権を表すleaseをテーブルに書く必要があります。
leaseとfencing tokenの違いは何ですか?
leaseは『このworkerがこの時刻までメッセージを所有する』という期限付きの所有権で、workerが落ちても期限が来れば他のworkerが引き継げます。fencing tokenはclaimのたびに増える番号で、確定時にこの番号が一致するworkerだけが書き込めます。leaseが失効した後に古いworkerが戻ってきても、番号が古いので書き込みが弾かれます。
fencing tokenがあれば外部APIの二重呼び出しも防げますか?
防げません。fencing tokenはデータベースへの確定書き込みを守りますが、古いworkerがleaseの失効前後に外部APIをすでに呼んでいれば、その副作用は残ります。外部側でも防ぐには、外部システムがtokenや冪等キーを検証する必要があります。Outboxでは、メッセージIDを冪等キーとして下流で重複排除する設計を併用します。
PostgreSQLをジョブキューとして使ってよいですか?
業務トランザクションと同じDBにメッセージを書けることがOutboxの価値なので、Outboxの配送キューとしてPostgreSQLを使うのは合理的です。ただし毎秒数万件のような高スループットや、長時間のジョブの大量保持には専用のキューやブローカーが向きます。queue depthと最古メッセージの年齢を監視し、DBの負荷が問題になったら発行先をブローカーに移す判断をします。

参考文献

友田

友田 陽大

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

この記事の実装を、案件として承ります

ORM選定・データモデル設計・ゼロダウンタイム移行を、設計から実装まで承ります

「どのORMを選ぶか」より「そのデータモデルが5年後の変更に耐えるか」が本質です。Prisma / Drizzle / SQLAlchemy などの選定、正規化と非正規化の線引き、N+1 と接続プールの設計、そして稼働中のサービスを止めないスキーマ移行までを一貫して設計・実装します。決済プラットフォームで信頼性レイヤーを主導し、冪等性と整合性を設計して本番の二重課金ゼロを維持した経験から、壊れたときに気づける・戻せるデータ層をつくります。

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

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

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

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

あわせて読みたい