結論から書きます。 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 方式に切り替えます。
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文が保証すること:
- 同時に claim した worker 同士は、互いに違う行を取る。
SKIP LOCKEDがここで効く。検証では、200件を4つの worker が7件ずつ同時に claim し続け、取得した ID は重複なしでちょうど200件だった。 - COMMIT 後も「processing・所有者・期限・世代」がテーブルに残る。 他の worker の claim は
status = 'pending'か「lease 切れのprocessing」しか取らないので、lease 中の行には触れない。 - 落ちた worker の行は、lease が切れれば自動的に再取得される。 再取得のたびに
lease_versionが上がる。 - 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_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 が切れる)を作れます。
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(整合性チェック)の設計 で扱います。