DB更新とイベント配信を分離するtransactional outbox設計|SQLite・Python・冪等性

「DB更新とイベント配信を分離するtransactional outbox設計|SQLite・Python・冪等性」の内容を表す技術イラスト

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は同じ対象内の順序判定に使います。

fieldtype役割
event_idstring重複排除に使う一意ID
event_typestring処理とschemaの選択
schema_versionintegerevent契約の版
aggregate_idstring更新対象のID
aggregate_versioninteger対象ごとの更新順序
dataobject公開してよい業務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を失わず、重複しても同じ結果へ収束させられます。

参考情報

参考になったらシェアしてください
  • URLをコピーしました!
  • URLをコピーしました!

この記事を書いた人

機械設計・油圧・CAD・Python・AIなど、ものづくりに関わる技術を扱っています。工学知識を整理・構造化し、設計や自動化に再利用できる形へ変えていくことを目指しています。

目次