1. 現場に蔓延する「2行のデスコード」

システムが成長し、非同期ワーカーやマイクロサービス、イベント駆動アーキテクチャ(EDA)を導入し始めた現場で、ほぼ100%の確率で混入する「極めて危険なコード」があります。

それが、**「データベースへの保存」と「メッセージブローカーへのパブリッシュ」をひとつの関数内に並べて逐次実行するコード**です。

「DB保存のあとにキューへ送信」という甘い罠

典型的なバックエンドの実装を見てみましょう。言語はGoを例に取りますが、JavaでもTypeScriptでもPythonでも構造はまったく同じです。

// 現場で量産される最も危険な「Dual Write」コード
func (s *OrderService) CreateOrder(ctx context.Context, order *Order) error {
    // 1. RDBのトランザクションを開始して注文レコードを保存
    tx, err := s.db.BeginTx(ctx, nil)
    if err != nil {
        return err
    }
    defer tx.Rollback()

    if err := s.orderRepo.InsertOrder(ctx, tx, order); err != nil {
        return err
    }

    // DBコミット
    if err := tx.Commit(); err != nil {
        return err
    }

    // 2. 外部メッセージキュー(Kafka / SQS / RabbitMQ)へイベントを送信
    event := OrderCreatedEvent{OrderID: order.ID, TotalAmount: order.Amount}
    if err := s.messageProducer.Publish(ctx, "order.events", event); err != nil {
        // 【爆弾】DBはコミット済みだが、イベント送信に失敗!
        // ロールバックはもう効かない。エラーログを吐く以外に何ができるのか?
        s.logger.Error("Failed to publish event", "order_id", order.ID, "err", err)
        return err
    }

    return nil
}

このコードをレビューしたとき、ジュニアや中級エンジニアは「綺麗に分離されていますね」「DBコミット後にイベントを投げているので安心です」と承認してしまいがちです。ローカル環境やCI環境のテストでも、何のエラーも出ずに快調にグリーンを灯すでしょう。

しかし、この2行の行間には、本番運用でシステムを確実に死に至らしめる断絶(Fault Tolerance Boundary)がぽっかりと口を開けています。

順序を入れ替えても地獄は終わらない

「DBコミットとキュー送信」という2つの独立したリモートシステムへの書き込み(これを分散システム論で **Dual Write(二重書き込み)** と呼びます)を行う以上、実行順序をどちらに倒しても破綻します。

パターンA: DBコミット → キュー送信(上記のコード)
DBのコミットが完了した直後、ネットワーク瞬断でキューへの通信がタイムアウトするか、あるいはデプロイやOOMキラーでPod/プロセスが強制停止(SIGKILL)されたらどうなるか?
結果: DBには注文が存在するが、メッセージキューにはイベントが一切流れない。下流の「決済サービス」や「出荷サービス」「在庫引き当てワーカー」は注文の存在すら認知できず、ユーザーは購入完了画面を見たのに商品が永遠に届かない「サイレント欠落事故」となる。
パターンB: キュー送信 → DBコミット
「じゃあ先にキューへ送信して、成功したらDBをコミットすればいいのでは?」と考える人がいます。しかしそれは更に悲惨な事故を生みます。
キューに送信した直後、DBのUNIQUE制約違反やデッドロックで `tx.Commit()` が失敗したらどうなるか?
結果: DBには注文が1件も存在しないのに、キューには「注文作成イベント」が先行して放流される。下流サービスがイベントを受け取り、存在しない幽霊注文に対して決済を実行し、二重請求やインベントリ不整合を引き起こす。

「失敗したらリトライすればいい」と安易に考えてはいけません。リトライ中にプロセスが死ねばイベントは失われ、逆にリトライが成功しても「実は1回目のリクエストがキュー側には届いていた」場合は重複配信が発生します。「2つの別々のデータストアに、同時に、不可分(Atomic)に書き込む手段は、通常のネットワーク越し通信には存在しない」——これが分散コンピューティングの残酷な真理です。

2. なぜ「分散トランザクション(2相コミット)」は現実解にならないのか?

コンピュータサイエンスを学んだ方なら、「2相コミット(2PC: Two-Phase Commit)やXAトランザクションを使えば、DBとメッセージブローカーを同一トランザクションで束ねられるのではないか?」と思い当たるかもしれません。

結論から言えば、現代のクラウドネイティブ環境において、2相コミットは運用保守の悪夢であり、真っ先に避けるべきアンチパターンです。

2PCが本番で忌避される3大理由

  1. コーディネーター障害による全体ロック死: 2PCは「準備フェーズ(Prepare)」と「コミットフェーズ(Commit)」の2段階を踏みますが、トランザクションコーディネーターがフェーズ途中で死んだ場合、関係するすべてのストレージがロックを保持したまま待機し続け、システム全体のリソースが枯渇して沈没します。
  2. スループットの破滅的低下: ネットワークラウンドトリップが複数回発生し、DBの行ロックを極めて長い時間保持するため、システムの同時実行性能が1/10〜1/100以下に激減します。
  3. クラウドサービス側の非対応: AWS SQS、Apache Kafka、GCP Pub/Subなどの近代的メッセージング基盤は、そもそも2PC(XAプロトコル)をサポートしていません。クラウドアーキテクチャの前提と根本的に噛み合わないのです。

同期的な分散トランザクションで整合性を担保しようとするアプローチは、マイクロサービスや非同期アーキテクチャが本来目指していた「システムの疎結合性」と「可用性」を自ら放棄する自殺行為に他なりません。

3. 救世主「Transactional Outboxパターン」の構造美

分散トランザクションを捨て、外部キューとのDual Writeを完全に防ぎながら、「100%確実にイベントを届ける」にはどうすればよいのか?

その答えが、Chris Richardson氏が提唱した名著『Microservices Patterns』の中核をなす設計思想、**「Transactional Outbox パターン(トランザクショナル・アウトボックス)」**です。

原理: 異種ストレージではなく「同一DBのACIDトランザクション」に閉じ込める

発想の転換は驚くほどシンプルで美しいものです。

「DB更新」と「外部キューへの送信」という2つの異なる仕組みに同時に書き込むから破綻するのです。ならば、「送信したいメッセージ(イベント)」を外部キューではなく、同じDBの中にある『outbox(送信箱)』テーブルにレコードとして保存してしまえばいいのです。

【従来のDual Write(アンチパターン)】
[App] ──(1. DB Commit)──> [PostgreSQL (orders)]
  │
  └───(2. Publish)──────> [Message Queue] 💥ここで通信障害やクラッシュが起きると不整合!

【Transactional Outbox パターン(堅牢設計)】
[App] ──(単一のローカルトランザクション)──┐
  │                                     │
  ├─(INSERT)─> [orders テーブル]        │(ACID特性により 100% 両方成功するか
  └─(INSERT)─> [outbox テーブル]        │  100% 両方ロールバックするかの二者択一)
                                        ┘
          │ (非同期で安全にリレー)
          ▼
[Outbox Relay Worker] ────(Publish)────> [Message Queue]

ビジネスエンティティ(orders)の変更と、イベント(outbox)の永続化は、**単一のRDBが数十年磨き上げてきたローカルACIDトランザクション(BEGIN 〜 COMMIT)**の保護下に置かれます。データベースの電源が落ちようがプロセスが吹き飛ぼうが、両方が保存されるか、両方が消えるかのどちらかしか起こり得ません。

そして、一旦DBに永続化されたOutboxメッセージを、別の軽量なバックグラウンドワーカーが安全にメッセージキューへ配送(Relay)します。

Outboxテーブルの最小スキーマ設計

Outboxパターンを実現するためのテーブル設計は、過度な装飾を排した極めてミニマルなもので十分です。PostgreSQLでの定義例を見てみましょう。

-- トランザクショナルOutboxテーブル
CREATE TABLE outbox_events (
    id            UUID PRIMARY KEY,
    aggregate_type VARCHAR(64) NOT NULL,  -- 集約名(例: 'order', 'user')
    aggregate_id   VARCHAR(64) NOT NULL,  -- 集約ID(例: order_id)
    event_type     VARCHAR(64) NOT NULL,  -- イベント名(例: 'OrderCreated')
    payload        JSONB NOT NULL,        -- イベントの本体データ
    created_at     TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP,
    processed_at   TIMESTAMPTZ            -- 配送完了日時(NULLなら未配送)
);

-- 未処理イベントを高速に吸い出すための部分インデックス(Partial Index)
CREATE INDEX idx_outbox_unprocessed 
ON outbox_events (created_at) 
WHERE processed_at IS NULL;

未処理メッセージの検索を爆速化するために、WHERE processed_at IS NULL という部分インデックス(Partial Index)を貼るのが運用上の最大のキモです。処理済みレコードが数百万件に膨らんでも、インデックスサイズは未処理の数十〜数百件分しか肥大化せず、スキャンコストは常に最小限に抑えられます。

4. Go言語とPostgreSQLで組む「ゼロ外部インフラ」完全実装

Outboxパターンを現場に導入しようとすると、一部のアーキテクトは「Debezium(CDC: Change Data Capture)を入れてPostgreSQLのWAL(Write-Ahead Log)をKafkaに直結しよう」といった巨大なインフラ構成を提案しがちです。

しかし、巨大なCDC基盤はKafkaクラスタやKafka Connectの維持管理コストを跳ね上げます。一般的なWebサービスや中規模マイクロサービスであれば、外部ミドルウェアを一切増やさず、PostgreSQLの FOR UPDATE SKIP LOCKED を使った数十行のGoワーカーだけで完璧にスケールします。

ステップ1: トランザクション内でイベントを同居させる

まずは書き込み側の実装です。外部キューのクライアントは登場せず、同一トランザクション内で outbox_events テーブルにレコードを流し込みます。

// ビジネスロジックとイベント永続化を単一トランザクションで実行
func (s *OrderService) CreateOrder(ctx context.Context, order *Order) error {
    tx, err := s.db.BeginTx(ctx, nil)
    if err != nil {
        return fmt.Errorf("begin tx: %w", err)
    }
    defer tx.Rollback()

    // 1. ordersテーブルへの保存
    if err := s.orderRepo.Insert(ctx, tx, order); err != nil {
        return fmt.Errorf("insert order: %w", err)
    }

    // 2. イベントペイロードのシリアライズ
    eventPayload, err := json.Marshal(OrderCreatedEvent{
        OrderID:     order.ID,
        UserID:      order.UserID,
        TotalAmount: order.TotalAmount,
    })
    if err != nil {
        return fmt.Errorf("marshal event: %w", err)
    }

    // 3. outboxテーブルへ同時書き込み
    const insertOutboxQuery = `
        INSERT INTO outbox_events (id, aggregate_type, aggregate_id, event_type, payload)
        VALUES ($1, $2, $3, $4, $5)
    `
    // IDには時間順序性を持つUUIDv7などを推奨
    eventID := uuid.New()
    if _, err := tx.ExecContext(ctx, insertOutboxQuery, 
        eventID, "order", order.ID, "OrderCreated", eventPayload,
    ); err != nil {
        return fmt.Errorf("insert outbox: %w", err)
    }

    // 両方が不可分に確定(Commit)
    return tx.Commit()
}

この時点で、ビジネスデータの整合性は完璧に守られました。ネットワークが途絶しようがDBが死のうが、不整合は1ミリも発生しません。

ステップ2: `SKIP LOCKED` で並列ワーカーの競合を消し去る

次に、DBに溜まった未処理イベントをキューへ中継する「Outbox Relay ワーカー」の実装です。ここで複数台のサーバー・Podが同時に動いてもロック競合を起こさない魔法の構文が、SQL標準(PostgreSQL / MySQL 8.0+)の FOR UPDATE SKIP LOCKED です。

// バックグラウンドで定期実行されるOutboxリレー処理
func (r *OutboxRelay) ProcessPendingEvents(ctx context.Context) error {
    tx, err := r.db.BeginTx(ctx, nil)
    if err != nil {
        return err
    }
    defer tx.Rollback()

    // 他のワーカーが処理中のレコードをスキップし、未処理イベントを最大50件ロック取得
    const selectQuery = `
        SELECT id, event_type, payload
        FROM outbox_events
        WHERE processed_at IS NULL
        ORDER BY created_at ASC
        LIMIT 50
        FOR UPDATE SKIP LOCKED
    `
    rows, err := tx.QueryContext(ctx, selectQuery)
    if err != nil {
        return err
    }
    defer rows.Close()

    type OutboxRecord struct {
        ID        uuid.UUID
        EventType string
        Payload   []byte
    }
    var records []OutboxRecord

    for rows.Next() {
        var rec OutboxRecord
        if err := rows.Scan(&rec.ID, &rec.EventType, &rec.Payload); err != nil {
            return err
        }
        records = append(records, rec)
    }

    if len(records) == 0 {
        return nil // 未処理イベントなし
    }

    // メッセージキューへ安全に送信
    for _, rec := range records {
        if err := r.publisher.Publish(ctx, rec.EventType, rec.Payload); err != nil {
            // 送信失敗時はコミットせずロールバック(次回のリトライ対象となる)
            return fmt.Errorf("publish failed: %w", err)
        }

        // 配送完了のマーク
        const markProcessedQuery = `
            UPDATE outbox_events 
            SET processed_at = CURRENT_TIMESTAMP 
            WHERE id = $1
        `
        if _, err := tx.ExecContext(ctx, markProcessedQuery, rec.ID); err != nil {
            return err
        }
    }

    return tx.Commit()
}

FOR UPDATE SKIP LOCKED を使うことで、ワーカーAが処理している行をワーカーBがブロックされることなく華麗にスルーし、次の行を並行して掴みます。余計なRedis分散ロックや複雑なオーケストレーションツールを使わずとも、**DBの機能だけで完全な並列分散キューワーカーが成立する**のです。

ステップ3: 忘れてはならない「受信側の冪等性(Idempotent Consumer)」

Transactional Outboxパターンを運用する上で、シニアエンジニアが絶対に肝に銘じておかなければならない鉄則があります。

それは、このパターンが保証するのは **「At-least-once(最低1回配送)」であり、「Exactly-once(厳密に1回)」ではない** という冷徹な現実です。

【重複配信が発生する不可避のシナリオ】
1. ワーカーが外部キュー(SQS)へPublishに成功
2. DBに「processed_at = NOW()」を書き込もうとした瞬間にネットワーク切断またはPod強制停止
3. DB上は「processed_at IS NULL」のまま残る
4. 再起動したワーカーが、同じメッセージを再度キューへPublishしてしまう!

分散システムの定理として、ネットワーク通信を介する以上、重複送信を100%ゼロにすることは不可能です。したがって、イベントを受信する側(Consumer)は、必ず同一イベントを2回以上受信しても安全な「冪等性(Idempotency)」を備えていなければなりません。

-- コンシューマー側の冪等性担保テーブル
CREATE TABLE processed_messages (
    message_id   UUID PRIMARY KEY,     -- 送信元OutboxのID
    processed_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP
);

コンシューマーは処理を開始する際、自身のローカルトランザクション内で processed_messages テーブルに対象の message_id をINSERTし、もし一意制約違反(Duplicate Key)が発生したら「すでに処理済み」と判定して何事もなくACKを返す。この**「送信側Outbox(At-least-once) + 受信側Idempotent Consumer(冪等処理)」**が組み合わさることで、システム全体として真の整合性が完結します。

5. まとめ: 「非同期」を語るなら、まず足元のACIDを信じよ

現代のソフトウェア開発では、「マイクロサービス」「イベント駆動」「非同期処理」といったバズワードが先行し、その華やかさの裏にある分散システムの代償——ネットワーク分断、二重書き込み、データロストの危険性——が軽視されがちです。

「DB更新のすぐ下にキュー送信を書く」という素朴なコードは、開発環境では完璧に動いているように見えて、本番のトラフィックと障害の波に晒された瞬間に、原因究明不能なデータ不整合を引き起こしてエンジニアを深夜の緊急障害対応で疲弊させます。

未知の外部ミドルウェアを無駄に増やしてインフラ構成を複雑怪奇にする前に、私たちが毎日使っているリレーショナルデータベースが備える強大な武器「ACID」をもう一度正しく信じること。それこそが、泥臭い本番運用を生き抜くシニアエンジニアの実践知です。