背景となるバックエンド面接の論点は、バックエンドエンジニアの面接対策AI:技術面接を突破する実践ガイド にもあります。本稿はその中の「メッセージキューをいつ使うか」を、Transactional Outbox の実装と障害境界に絞って新たに解説します。
結論
注文を保存した直後にメッセージブローカーへ publish する設計では、DB 更新とイベント公開のどちらかだけが成功する時間帯をなくせません。Transactional Outbox は、業務データと「これから公開すべきイベント」を同じ DB トランザクションに入れ、別ワーカーが後で配送することで、この不整合を回復可能にするパターンです。
ただし、これで「exactly once」になるわけではありません。実運用の答えは次の 3 点です。
- DB への書き込みと outbox 行の追加を同じトランザクションにする
- worker は lease 付きでイベントを取得し、少なくとも 1 回配送する
- consumer は
eventIdを使って冪等に処理する
システム設計面接なら、ここまでを「失敗時に何が残り、次に誰が回収するか」で説明できると強いです。
なぜ「保存してから publish」が壊れるのか
次の素直なコードには、再試行だけでは埋められない穴があります。
await db.insertOrder(order);
await broker.publish({ type: "order.created", orderId: order.id });
insertOrder の後にプロセスが停止すれば、注文は存在するのに downstream はそれを知りません。順番を逆にすると、イベントだけが届いて注文が見えない可能性があります。DB と broker は別の耐障害ドメインなので、アプリケーションだけで原子的に両方を commit することはできません。
そこで、最初の commit の対象を DB だけに限定します。
BEGIN;
INSERT INTO orders (id, customer_id, status)
VALUES (:orderId, :customerId, 'created');
INSERT INTO outbox (id, topic, payload, status, created_at)
VALUES (:eventId, 'order.created', :payload, 'pending', now());
COMMIT;
commit に成功した時点で、注文と配送待ちイベントは必ず対になって残ります。broker が落ちていても、worker が後で outbox を読み直せます。
worker が守るべき境界
worker の正しい順序は、次の通りです。
-
pending、または lease が期限切れのイベントを 1 件 claim する - broker へ
eventIdを含めて publish する - 成功後に outbox を
deliveredにする
最も重要なのは 2 と 3 の間です。publish の成功直後に worker が停止すると、次の worker は同じイベントを再送します。これはバグではなく、at-least-once delivery の仕様です。重複を consumer 側で止める前提を明示します。
PostgreSQL なら複数 worker の claim を FOR UPDATE SKIP LOCKED で行えます。処理中フラグだけでは worker のクラッシュで永久に詰まるため、locked_until を持つ lease にします。
WITH picked AS (
SELECT id
FROM outbox
WHERE status = 'pending'
OR (status = 'processing' AND locked_until < now())
ORDER BY created_at
FOR UPDATE SKIP LOCKED
LIMIT 1
)
UPDATE outbox
SET status = 'processing',
locked_until = now() + interval '30 seconds',
attempts = attempts + 1
WHERE id IN (SELECT id FROM picked)
RETURNING *;
小さく実行できる TypeScript 例
本番では上の SQL を DB に対して実行します。以下は、同じ状態遷移をメモリ上で検証する最小例です。Bun でそのままテストできます。
import { expect, test } from "bun:test";
type Status = "pending" | "processing" | "delivered";
type Event = {
id: string;
orderId: string;
status: Status;
lockedUntil: number;
};
class OutboxDb {
private orders = new Set<string>();
private events = new Map<string, Event>();
transaction(fn: (tx: {
createOrder: (orderId: string, eventId: string) => void;
}) => void) {
// commit 前のコピーにだけ書く。例外なら元の状態へ触れない。
const nextOrders = new Set(this.orders);
const nextEvents = new Map(
[...this.events].map(([id, event]) => [id, { ...event }]),
);
fn({
createOrder: (orderId, eventId) => {
if (nextOrders.has(orderId) || nextEvents.has(eventId)) {
throw new Error("duplicate id");
}
nextOrders.add(orderId);
nextEvents.set(eventId, {
id: eventId,
orderId,
status: "pending",
lockedUntil: 0,
});
},
});
this.orders = nextOrders;
this.events = nextEvents;
}
claim(now: number): Event | undefined {
const event = [...this.events.values()].find(
(item) =>
item.status === "pending" ||
(item.status === "processing" && item.lockedUntil <= now),
);
if (!event) return;
event.status = "processing";
event.lockedUntil = now + 30_000;
return { ...event };
}
markDelivered(eventId: string) {
const event = this.events.get(eventId);
if (!event) throw new Error("event not found");
event.status = "delivered";
}
statusOf(eventId: string) {
return this.events.get(eventId)?.status;
}
}
class IdempotentConsumer {
private processed = new Set<string>();
readonly orders: string[] = [];
consume(event: Event) {
if (this.processed.has(event.id)) return;
this.processed.add(event.id);
this.orders.push(event.orderId);
}
}
test("publish 後に worker が停止しても、再配送で業務処理は重複しない", () => {
const db = new OutboxDb();
const consumer = new IdempotentConsumer();
db.transaction((tx) => tx.createOrder("order-1", "event-1"));
const first = db.claim(0)!;
consumer.consume(first); // broker への publish は成功
// ここで worker が停止し、markDelivered は実行されなかったとする
const retry = db.claim(30_001)!; // lease 切れを別 worker が回収
consumer.consume(retry);
db.markDelivered(retry.id);
expect(consumer.orders).toEqual(["order-1"]);
expect(db.statusOf("event-1")).toBe("delivered");
});
test("業務データと outbox は例外時に一緒に rollback される", () => {
const db = new OutboxDb();
expect(() =>
db.transaction((tx) => {
tx.createOrder("order-2", "event-2");
throw new Error("rollback");
}),
).toThrow("rollback");
expect(db.claim(0)).toBeUndefined();
});
この例で再配送が安全なのは、consumer が event.id を保存しているからです。実際には consumer の DB に processed_events(event_id primary key) を置き、業務更新と同じトランザクションで INSERT します。重複した event_id の INSERT が一意制約違反なら、すでに処理済みとして終了します。
面接で確認したいトレードオフ
1. なぜ分散トランザクションを使わないのか
2PC は coordinator、参加者のロック、障害時の可用性低下を持ち込みます。DB に確実に残せる outbox と冪等 consumer の組み合わせは、可用性を保ちながら最終的な整合性を得る実装です。
2. outbox は増え続けないか
delivered 行を監査保持期間の後に削除、またはアーカイブします。ただし削除は consumer の重複排除保持期間と矛盾させません。少なくとも、遅延再送が起こり得る期間は eventId の追跡情報を残します。
3. 順序が必要ならどうするか
単一集約(たとえば 1 注文)内の順序が必要なら、aggregate_id と連番をイベントに持たせます。全注文の完全なグローバル順序を一台の worker で守ると throughput を失うため、必要な順序の範囲を要件として先に確認します。
4. 監視する値は何か
最低限、pending 件数、最古イベントの滞留時間、lease 切れ再試行数、配送失敗数を監視します。特に最古イベントの age は「注文が作られてから downstream に届くまで」の遅延を直接示します。
まとめ
Transactional Outbox の本質は「broker を失敗しないものにする」ことではありません。失敗しても DB に回復の手掛かりを残し、重複しても consumer の結果を 1 回にすることです。
設計を説明するときは、次の一文に戻ると整理しやすくなります。
DB commit 後に残る outbox を再送し、publish と完了更新の間の停止は
eventIdによる冪等 consumer で吸収する。