DBの更新後にWebhookやmessage brokerへ通知すると、「DB更新は成功したが通知は失敗した」「通知は届いたがDB更新はrollbackした」という食い違いが起こり得ます。DBと外部systemを同じtransactionへ入れられないことが原因です。
transactional outboxは、業務dataの変更と「後で送るevent」を同じDB transactionへ保存し、外部配信を別processへ分離するpatternです。outbox tableの構造、SQLiteとPythonによる実装、再送・重複・順序・秘密情報の扱いを整理します。
二重書き込みが生む不整合
部品masterを更新し、別systemへasset.updated eventを送る処理を考えます。
DBの部品recordを更新
↓
外部APIまたはmessage brokerへeventを送信
| 結果 | 問題 |
|---|---|
| DB成功・送信失敗 | 受信側が更新を知らない |
| 送信成功・DB失敗 | 存在しない更新を受信側が処理する |
| 送信成功・応答消失 | 再送でeventが重複する |
AWS Prescriptive Guidanceも、この二重書き込みの片方だけが失敗するとsystem間でdataが不整合になると説明しています。HTTP送信をDB transactionの途中へ置いても、外部systemが受理した処理をDB rollbackで取り消すことはできません。
transactional outboxの仕組み
業務recordとevent recordを同じDBへ保存します。
API request
↓
同一DB transaction
├─ 業務tableを更新
└─ outbox tableへeventをINSERT
↓ commit
relayが未配信eventを取得
↓
外部systemへ送信
↓
配信済みstatusを記録
業務更新とoutbox INSERTのどちらかが失敗すれば両方をrollbackします。commit後にrelayが停止しても、event recordは残るため再開後に送れます。
ただし、外部送信後、配信済み更新の前にrelayが停止すると同じeventを再送します。outboxは通常、少なくとも1回の配信と受信側の冪等性を組み合わせます。「必ず1回だけ」を安易に保証してはいけません。
event envelopeを決める
eventは識別子、対象、種類、版、発生日時、payloadを持つobjectにします。
{
"event_id": "evt-01K6W3P8Y7",
"event_type": "asset.updated",
"schema_version": 1,
"aggregate_type": "asset",
"aggregate_id": "asset-042",
"aggregate_version": 8,
"occurred_at": "2026-10-07T01:20:00Z",
"data": {
"name": "Hydraulic unit B",
"material": "S45C"
}
}
event_idは配信試行ではなく、業務eventそのものの不変IDです。再送でも変更しません。aggregate_idは更新対象、aggregate_versionは同じ対象内の順序判定に使います。
| field | type | 役割 |
|---|---|---|
event_id | string | 重複排除に使う一意ID |
event_type | string | 処理とschemaの選択 |
schema_version | integer | event契約の版 |
aggregate_id | string | 更新対象のID |
aggregate_version | integer | 対象ごとの更新順序 |
data | object | 公開してよい業務data |
日時だけで順序を決めると、同時刻や時計差の問題があります。順序が重要なら単調増加versionを持たせます。
SQLiteにoutbox tableを作る
業務tableとoutbox tableを同じdatabaseへ置きます。
CREATE TABLE assets (
asset_id TEXT PRIMARY KEY NOT NULL,
name TEXT NOT NULL,
material TEXT NOT NULL,
version INTEGER NOT NULL CHECK (version > 0),
updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP
);
CREATE TABLE outbox_events (
event_id TEXT PRIMARY KEY NOT NULL,
aggregate_type TEXT NOT NULL,
aggregate_id TEXT NOT NULL,
aggregate_version INTEGER NOT NULL CHECK (aggregate_version > 0),
event_type TEXT NOT NULL,
schema_version INTEGER NOT NULL CHECK (schema_version > 0),
payload_json TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'pending'
CHECK (status IN ('pending', 'published', 'failed')),
attempt_count INTEGER NOT NULL DEFAULT 0,
available_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
published_at TEXT,
last_error TEXT,
UNIQUE (aggregate_type, aggregate_id, aggregate_version)
);
CREATE INDEX idx_outbox_delivery
ON outbox_events (status, available_at, created_at);
payload_jsonはJSON textとして保存しますが、保存前に型、必須field、NaNなどを検証します。配信後すぐ削除すると調査しにくいため、まずpublishedとして保持し、保持期限後にarchiveまたは削除します。
Pythonで更新とeventを同時に保存する
次の例は、assetのversionが期待値と一致するときだけ更新し、同じtransactionでoutbox eventを追加します。
import json
import sqlite3
from datetime import datetime, timezone
def update_asset(
conn: sqlite3.Connection,
*,
event_id: str,
asset_id: str,
expected_version: int,
name: str,
material: str,
) -> None:
new_version = expected_version + 1
event = {
"event_id": event_id,
"event_type": "asset.updated",
"schema_version": 1,
"aggregate_type": "asset",
"aggregate_id": asset_id,
"aggregate_version": new_version,
"occurred_at": datetime.now(timezone.utc)
.isoformat()
.replace("+00:00", "Z"),
"data": {"name": name, "material": material},
}
payload = json.dumps(
event,
ensure_ascii=False,
separators=(",", ":"),
allow_nan=False,
)
with conn:
cursor = conn.execute(
"""
UPDATE assets
SET name = ?, material = ?, version = ?,
updated_at = CURRENT_TIMESTAMP
WHERE asset_id = ? AND version = ?
""",
(name, material, new_version, asset_id, expected_version),
)
if cursor.rowcount != 1:
raise ValueError("対象がないか、versionが競合しています")
conn.execute(
"""
INSERT INTO outbox_events (
event_id, aggregate_type, aggregate_id,
aggregate_version, event_type,
schema_version, payload_json
) VALUES (?, ?, ?, ?, ?, ?, ?)
""",
(
event_id, "asset", asset_id, new_version,
"asset.updated", 1, payload,
),
)
with conn:内で例外が起きればrollbackします。event IDやaggregate versionの重複でINSERTが失敗した場合も、asset更新だけが残りません。SQLiteのtransaction仕様では、transactionはCOMMITまたはROLLBACKまで継続します。
同じAPI requestが再送されるなら、clientのidempotency keyとevent_idの対応を保存します。同じkeyへ毎回新しいevent IDを発行してはいけません。
relayは再送を前提にする
relayは未配信recordを小さなbatchで読み、外部systemへ送ります。処理の要点は次のとおりです。
def relay_once(conn, publish) -> bool:
row = conn.execute(
"""
SELECT event_id, payload_json
FROM outbox_events
WHERE status = 'pending'
AND available_at <= CURRENT_TIMESTAMP
ORDER BY created_at, event_id
LIMIT 1
"""
).fetchone()
if row is None:
return False
event_id, payload = row
try:
publish(event_id, payload)
except Exception as exc:
with conn:
conn.execute(
"""
UPDATE outbox_events
SET attempt_count = attempt_count + 1,
available_at = datetime('now', '+1 minute'),
last_error = ?
WHERE event_id = ? AND status = 'pending'
""",
(str(exc)[:500], event_id),
)
return False
with conn:
conn.execute(
"""
UPDATE outbox_events
SET status = 'published',
attempt_count = attempt_count + 1,
published_at = CURRENT_TIMESTAMP,
last_error = NULL
WHERE event_id = ? AND status = 'pending'
""",
(event_id,),
)
return True
この最小例はrelayが1processだけ動く前提です。複数processでは同じrecordを同時取得しないよう、利用DBに応じてatomicなclaim、lease期限、row lockを設計します。
通信timeoutを未送信と断定してはいけません。相手が受理した後にresponseだけ失われた可能性があります。publish()にはevent IDをmessage keyやidempotency keyとして渡します。
受信側もevent IDで重複排除する
受信側はevent IDを一意制約で記録し、業務処理と同じtransactionで重複を判定します。
CREATE TABLE processed_events (
event_id TEXT PRIMARY KEY NOT NULL,
processed_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP
);
def consume_once(conn, event: dict) -> bool:
with conn:
cursor = conn.execute(
"""
INSERT INTO processed_events (event_id)
VALUES (?)
ON CONFLICT (event_id) DO NOTHING
""",
(event["event_id"],),
)
if cursor.rowcount == 0:
return False
apply_asset_update(conn, event)
return True
業務処理が失敗すればprocessed IDもrollbackされ、修正後に再処理できます。IDだけ先にcommitすると、業務処理に失敗しても次回を重複扱いしてしまいます。
順序・schema・異常系を決める
event IDだけでは順序逆転を解決できません。aggregate_versionが7の後に6が届いた場合、古いeventで最新版を上書きしないようにします。
- 現在versionより小さい:反映済みとして無視
- 現在versionより1大きい:正常に適用
- 2以上大きい:欠落の可能性として保留、または正本APIから再取得
schema_versionが未対応、JSONが不正、必須fieldが欠落しているeventは、自動再試行では直りません。dead-letter領域へ隔離します。timeoutや一時的な5xxはbackoff後に再試行し、認証失敗はcredentialを直すまで停止します。
pollingとCDCを使い分ける
小規模systemはoutbox tableのpollingから始められます。規模が大きい場合は、DBのchange data captureでINSERTを検出する方法があります。
Debezium Outbox Event Routerは、outbox tableの変更をcaptureし、一意event ID、aggregate ID、type、payloadをmessageへ変換します。ただしCDCを使っても、受信側の冪等性やschema管理が不要になるわけではありません。
セキュリティとCAD・BOMへの展開
outboxは再送と監査のためにpayloadを保持します。API key、token、passwordを含めず、個人情報や機密属性も必要最小限にします。relayのcredentialはsecret管理機能から取得し、last_errorやlogへresponse bodyを無制限に保存しません。
大きなCAD fileや帳票をpayloadへ埋め込まず、file ID、revision、content hashを送り、認可済みAPIから取得する構成が扱いやすくなります。図面承認、BOM更新、計算完了、部品master変更なども、共通envelopeとrelayへ載せられます。
まとめ
transactional outboxは、同じDB transactionで「送るべきevent」を確定し、外部配信を分離する設計です。
- 業務更新とoutbox INSERTを同じtransactionでcommitする
- relayはcommit済みrecordだけを送る
- event IDを再送でも変えない
- 受信側を冪等にする
- aggregate versionで順序逆転と欠落を検出する
- schema version、再試行、dead-letter、保持期限を決める
- pollingとCDCは規模と運用要件で選ぶ
- payloadからsecretと不要な機密dataを分離する
この構造なら、API、queue、Webhook、DB、CAD、BOMの連携で障害が起きても、未配信eventを失わず、重複しても同じ結果へ収束させられます。

