イベントを失わないイベント駆動マイクロサービス:Transactional Outbox の実践
Cloud Run と Spanner の上で、注文・決済・下流サービスの整合性をどう守ったか——Transactional Outbox パターンとタスクキューのリレーワーカー、そして今ならこう作るという反省。
私が一番嫌いな本番バグのクラスは、何もクラッシュさせません。決済はキャプチャされたのに注文は凍ったまま、アラームは鳴らない——システムが静かに自分自身と食い違っていくのです。そして、この種のバグを追いかけると、ほぼ毎回同じ2行の無邪気なコードに行き着きます:行をコミットし、それからイベントをパブリッシュする。プロセスがこの2行の間で死ねば、下流サービスは注文のことを永遠に知りません。順序を逆にすれば、ロールバックによって「存在しない注文」の知らせが飛びます。
私が関わったコンシューマー向けECプラットフォーム——GoのマイクロサービスをCloud Runで動かし、Cloud Spannerを主データベースに、十数のサービスが決済ゲートウェイを含む複数の下流システムと連携して複数ステップの購入フローをオーケストレーションする構成——では、この二重書き込みこそ、お金の経路で絶対に許容できない故障モードでした。対応する注文遷移のない決済キャプチャ、発行されなかったイベントを待ち続ける注文。アラートではなく、突合クエリで発見される類のずれです。
解決策は地味で、文書化され尽くしていて、完全に有効でした:Transactional Outbox です。
核心のアイデア
メッセージバスへのパブリッシュを副作用として行うのではなく、状態変更と同じデータベーストランザクションの中にイベントを書き込みます:
-- 1つの Spanner read-write トランザクション
INSERT INTO Orders (OrderId, Status, ...) VALUES (@id, 'CONFIRMED', ...);
INSERT INTO OutboxEvents (EventId, Topic, Payload, CreatedAt, PublishedAt)
VALUES (@eventId, 'order.confirmed', @payload, PENDING_COMMIT_TIMESTAMP(), NULL);
両方の行が存在するか、どちらも存在しないか。すでに対価を払っているデータベースの原子性という保証が、メッセージングまで覆うようになります。
あとは独立したリレーワーカーが未発行の行を読み、外へ押し出します。配送機構にはCloud Tasksを使いました:リレーはイベントごとにコンシューマーのエンドポイントを宛先とするタスクをエンキューし、行を発行済みにマークし、リトライと指数バックオフはタスクキューに任せます。
func (r *Relay) Tick(ctx context.Context) error {
events, err := r.repo.FetchUnpublished(ctx, batchSize)
if err != nil {
return err
}
for _, ev := range events {
if err := r.tasks.Enqueue(ctx, ev); err != nil {
return err // 行はそのまま残す。次のTickでリトライされる
}
if err := r.repo.MarkPublished(ctx, ev.ID); err != nil {
return err // タスクが二重発火しうる——コンシューマー側で重複排除が必須
}
}
return nil
}
最後のエラーパスのコメントが契約のすべてです:Outboxが与えるのは at-least-once であって、exactly-once ではありません。 リレーがエンキューとマークの間でクラッシュすれば、イベントは2回飛びます。トランスポート層で exactly-once を追いかけるのは負け戦で、問題はコンシューマー側に押し出すべきものです。
冪等なコンシューマーがパターンの半分
at-least-once 配送は、再配送が無害であって初めて成立します。各ドメインはそれを自分の言葉で強制していました——リレーはエンキュー前に発行済みフラグを確認し、タスクキューの重複エラーを握りつぶす。決済経路はリクエストごとの冪等性キーで重複を弾く。そして一般化できる形、つまり真似する価値のある形は、状態変更と一緒にコミットされる重複排除レコードです:
// コンシューマー自身のトランザクションの中で
applied, err := s.repo.TryRecordEvent(ctx, ev.ID) // INSERT OR IGNORE 相当
if err != nil || !applied {
return err // 重複——ACKして先へ進む
}
// ... 同じトランザクション内で状態変更を適用する
重複排除レコードと状態変更が原子的にコミットされる。つまりコンシューマーは自分専用のミニチュア「逆向きOutbox」を持つわけです。コンシューマーがこの線を守ったところでは、重複は文字どおり「非イベント」になりました。
複数ステップのフローをオーケストレーションする
購入フローは複数のサービスにまたがります:在庫確保、決済オーソリ、注文確定、フルフィルメント起動、通知送信。私たちは意図的にサーガフレームワークに手を伸ばしませんでした。各ステップは「イベント → 冪等なハンドラー → 次のイベント」の連鎖で、どのリンクも静かに脱落しないことをOutboxが保証します。
これを管理可能に保った設計ルールは2つ:
- イベントは事実を運ぶ。命令ではない。
order.confirmedであってsend_email_pleaseではありません。事実が自分にとって何を意味するかはコンシューマーが決める。これによりプロデューサーは受け手を知らずに済み、プロデューサーに触れずにコンシューマーを追加できます。 - すべてのフローに単一のオーナーサービスを置く。 注文のステートマシンは1つのサービスが所有し、注文ライフサイクルのイベントを書けるのはそのサービスだけ。他は反応するのみ。深夜2時に何かがおかしくなったとき、「この注文は実際どの状態なのか」を見に行く場所が正確に1か所になります。
決済ゲートウェイについては特に、カード・ウォレット・コンビニ払いのすべての決済手段を、専用のゲートウェイ抽象化サービスの背後に隔離しました。外部の決済プロバイダーは独自のリトライ意味論、コールバック、失敗の語彙を持っています。それをコアの注文サービスに漏らせば、私たちのステートマシンがサードパーティの癖に結合してしまう。ゲートウェイサービスはプロバイダーのコールバックを、他のすべてが消費するのと同じクリーンな内部イベントに翻訳しました。
運用ノート
設計ドキュメントよりも本番で効いたこと:
- Outboxテーブルの成長。 発行済みの行は保持期間の後、スケジュールジョブが削除します。Spannerは大きなテーブルに強いとはいえ、無限に伸びるOutboxはリレーの
FetchUnpublishedスキャンを遅くします。(PublishedAt, CreatedAt)にインデックスを張り、ワーキングセットを小さく保つこと。 - 順序。 タスクキューは順序を保証しません。順序が重要な場面(ステートマシンの遷移)では、コンシューマーは到着順を信じるのではなく遷移自体を検証します——
order.confirmedより先に来たorder.shippedは適用せず、退避させてリトライ。 - 可観測性。 リレーは
oldest_unpublished_ageをメトリクスとしてエクスポートします。このゲージ1つでほぼすべての故障モードが捕まります:リレーの停止、タスクキューの詰まり、コンシューマーのエラー。アラートはキューの深さではなく、この「年齢」に張ること。 - コンテキストのキャンセル。 微妙な本番バグをひとつ:goroutineで発火したキャッシュ書き込みがリクエストのコンテキストを継承しており、ハンドラーが戻った瞬間に書き込み途中でキャンセルされました。修正は「賢さをやめる」こと——リクエストパス上で同期的に実行する(ポストモーテム全文)。
今ならどうするか
もう一度やり直すなら、Outboxの配管コードは生成します。初期はリレーとリポジトリのコードをサービスごとに手書きしていて、実装間のドリフトが摩擦の大半を生みました。後に共有の内部ライブラリとDAOコードジェネレーターに集約してからは、このパターンの採用コストはほぼゼロになりました。
そしてパターンの導入は、最初の整合性インシデントの前にやること。Outboxのコストは追加テーブル1つと小さなワーカー1つです。代替案のコストは、注文と決済のレコードをSQLスクリプトで突合する週末——そしてチェックアウトに対するユーザーの信頼です。
二重書き込みは、まだ踏んでいないだけのバグです。イベントはトランザクションの中に書きましょう。
この記事は同じプラットフォームを扱う短いシリーズの起点です:1つのインターフェース、4つの支払い方法、サービスメッシュなしのサービス間認証、goroutineが間違った道具になるとき。