2つのOutbox、1つのパターン

コマンドにはCloud Tasks、イベントにはPub/Sub——同じTransactional Outbox、同じリレー、同じHTTP配送。本当の分岐点は中間層と、ルーティングを誰が所有するか。

私は非同期アーキテクチャをひとつの問いで評価するようになりました:新しいサービスが何かに反応する必要が生じたとき、誰のTerraformが変わるのか? この問いを学んだプラットフォーム——Cloud Run上の十数のGoマイクロサービス、足元にCloud Spanner——では、その答えがメッセージング設計全体をきれいに二分していて、その分割はどんな本よりも「コマンド対イベント」を教えてくれました。

このシステムのすべての非同期ホップは同じ動きです:メッセージをビジネストランザクションの中でSpannerに書き、リレーがそれをマネージドなGCPサービスへ押し出し、認証付きHTTP POSTとして受け取る。最初のステップがなぜ重要かはTransactional Outboxの記事が扱っています。この記事は中間の分岐について:ポイントツーポイントのコマンドにはCloud Tasks、ファンアウトするイベントにはPub/Sub——そして興味深い違いはスループットや意味論ではなく、所有権だという話です。

共有されるパターン

サービスが直接パブリッシュすることはありません。ビジネスデータと同じSpannerトランザクションの中で、2つのOutboxテーブル——tasksMessages か pubsubMessages——のどちらかに行を挿入します。メッセージがデータなしに存在することも、その逆もあり得ません。テーブルごとに専用のシングルインスタンスCloud Runリレーが未発行の行をポーリングし、マネージドサービスへ渡します。そして受信側では、コンシューマーはキュークライアントのコードを一切走らせません:配送は常に、ロードバランサーの背後にある普通のルートへのOIDC署名付きHTTP POSTです。

この最後の性質は立ち止まる価値があります。どちらの経路もプレーンなHTTPで配送されるため、メッセージハンドラーはただのエンドポイントです——同期ルートと同じミドルウェア、同じ認証コンテキストの伝播、同じロギング、同じローカルテストの体験。ビジネスロジックが住む場所では、メッセージングインフラは見えないのです。

分岐点:ルーティングを誰が所有するか

Spannerの中では2種類のメッセージ行はほぼ同一に見えます。見分けるのはフィールド1つと、Terraformの置き場所です:

Cloud Tasks — コマンドPub/Sub — イベント
ルーティングフィールドqueueId → 事前に束縛された単一エンドポイントtopicId → ブロードキャストチャネル
宛先の定義場所プロデューサーのTerraform——キューとターゲットパスが一緒各コンシューマーのTerraform——自分のサブスクリプションと自分のパス
プロデューサーは受信者を知る?はい——HTTPヘッダーまで組み立てるいいえ——サブスクライバーの知識ゼロ
配送されるコピーちょうど1ハンドラーサブスクリプションごとに1つ
フロー制御キュー単位のディスパッチレート+並行数上限コンシューマー単位にはなし——プッシュはスケールアップする
重複排除タスクID = メッセージID。重複エンキューは AlreadyExists を返すなし——コンシューマー自身が冪等であること
受信者の追加プロデューサーのインフラを変更コンシューマーがサブスクリプションを追加。プロデューサーは無傷
意味論「後でこのエンドポイントを呼べ。リトライ付きで」「これが起きた。関心があれば反応せよ」

3行目を二度読んでください——それが区別のすべてです。コマンドはプロデューサーの用事です:誰が行動すべきかを正確に知っているので、宛先はプロデューサーの設定に属します。イベントはコンシューマーの用事です:プロデューサーは事実(user-account-created)に名前を付けるだけで、関心のあるサービスは自分のインフラで自分を接続します。プラットフォームの技術ポリシーは、すべての非同期処理のデフォルトをCloud Tasksとし、Pub/Subは複数サービスが反応すべき本物のドメインイベントの一握りのために予約しました——これにより「これはイベントであるべきか?」が、デフォルトではなく意図的な設計の会話であり続けたのです。

Path Aの実例:サービスが自分自身に仕事をキューイングする

最もきれいなコマンドの例は、プロデューサーとコンシューマーが同じサービスであるケースです。外部パートナーが注文完了の通知をPOSTすると、コマースサービスはビジネス行とタスク行を1つのトランザクションで挿入し、即座に200を返します。リレーが行を拾い、タスクをバッファし、Cloud Tasksがそれを——同じサービスの別のエンドポイントへ——POSTで送り返す。そこで重い仕事が実行されます:メールゲートウェイ経由の確認メール、課金サービス経由の決済キャプチャ。

パートナーのリクエストはデータベースのコミットで終わります。高価なものはすべてタスクに乗り、非2xxでのリトライとキュー単位のレート制限が付いてくる——他人のHTTPリクエストの中で走るべきではなかった仕事のための、耐久性です。

重複排除ルールが存在する理由:リレーのクラッシュウィンドウ

Cloud TasksのタスクIDをOutboxの messageId に設定することは、正確にこの隙間を閉じます:

  1. リレーが行を読み、Cloud Tasksにタスクを作成する。
  2. リレーが publishedAt をSpannerに書き戻す前にクラッシュする。
  3. 次のTickでは行はまだ未発行に見える——リレーは再びタスクを作成する。
  4. Cloud Tasksは同じタスクIDを見て AlreadyExists で拒否する。
  5. リレーはそれを成功として扱い、行を発行済みにマークする。

この1行——リレーが AlreadyExists を明示的に無視すること——がなければ、リレーがクラッシュするたびに注文メールの二重送信や決済の二重キャプチャが起こり得ました。Pub/Subには同等の仕組みがありません。だからこそポリシーのもう半分が存在します:Pub/Subのコンシューマーは自分自身で冪等でなければならない。

Path Bの実例:サインアップのファンアウト

ユーザーがサインアップすると、IDサービスはユーザー行とイベント行を1つのトランザクションで挿入し、トピック名だけを名指しします。リレーがパブリッシュし、プッシュサブスクリプションがそれぞれにコピーを配送します:ポイントサービスはアカウントを初期化し、通知サービスはデフォルト設定を作成し——そしてIDサービス自身もサブスクライブして、自分のハンドラーで紹介コードを割り当てます。

この最後のサブスクリプションが、この設計全体で私のお気に入りの細部です。プロデューサーが自分のトピックを購読していることは、デカップリングが本物である証明です:プロデューサー自身の後続処理でさえ、サインアップのトランザクションではなくイベントに乗る。サインアップは最小限の原子的書き込みのままで、下流のすべて——同一サービス内の下流を含めて——は反応です。後に新しいサービスがサインアップに反応する必要が生じたとき、変更はそのサービスのTerraformのサブスクリプション1つでした。IDサービスは何も知らないままです。

トランスポートは1つ、選択肢は除外

このシステムの運用感覚を形づくった最後の観察:gRPCはどこにもありません。 同期API、Cloud Tasksのディスパッチ、Pub/Subのプッシュ——サービス間のすべての呼び出しは、OIDCトークンと伝播される認証コンテキストヘッダーを伴うHTTP/JSONです。gRPCが現れるのはクラウドSDKの内部とローカルエミュレーターだけ。トランスポートが1つなら、ミドルウェアスタックも1つ、トレースの読み方も1つ、緊急時にハンドラーを curl する方法も1つ——そしてそれが、非同期経路でも「コンシューマーはただのエンドポイント」を真にするのです。

私が持ち歩く決定ルール:コマンドはプロデューサーが所有するインフラを通し、イベントはコンシューマーが所有するインフラを通す。 メッセージが表のどちら側に座るか決められないなら、あなたはまだ「プロデューサーは聞き手を知ってよいのか」を決めていません——そしてキュー技術ではなく、それこそが本当のアーキテクチャ上の選択なのです。