プログラマのための IT 教科書

非同期処理とメッセージング

この部の 6 / 13 章 ・ 全体で 48 / 76 章 ・ 読了目安 45 分

この章を読むとできるようになること
  • 重複配信を前提にコンシューマを書ける
  • DB とメッセージ送信の整合性を保てる
  • 非同期の障害を検知できる

マイクロサービスの実務では、サービス同士が同期的に呼び合う構成を見ました。 呼んで、待って、返ってくる。単純で分かりやすい方法です。

しかし、この構成には限界があります。

[注文サービス] → [在庫] → [決済] → [メール送信] → [分析]
                                        ↑
                    メール基盤が落ちていると、注文自体が失敗する

注文を受け付けることと、確認メールを送ることは、同じ重要度ではありません。 それなのに、同期的に繋いだ結果、メールの障害が注文を止めています。

そこでメッセージングを使います。

[注文サービス] --メッセージを置く--> [キュー] --取り出す--> [メール送信]
     ↓
  すぐ完了を返せる。メール基盤が落ちていても、注文は成立する

いつ使うか

使う場面理由
時間のかかる処理動画の変換、帳票の生成。待たせない
落ちてもよい処理通知、分析、キャッシュの更新
負荷の平準化一斉配信のような波を、処理側のペースで消化する
1つの出来事を複数が使う注文確定 → 在庫・通知・分析(ドメイン駆動設計の実践のドメインイベント)
外部サービスへの依存を切る相手の障害を、自分の障害にしない
使ってはいけない場面
□ 呼び出し元が結果をすぐ必要とする(残高照会、認証)
□ 厳密な順序が本質的に必要な処理
□ 単純な CRUD

非同期にすると、結果がいつ返るか分からなくなります。 「すぐ答えが要る」処理を非同期にすると、 状態をポーリングする仕組みが必要になり、かえって複雑になります。

非同期は複雑さと引き換えに、独立性を買う手段です。

キューとトピック

キュー(1対1)      1つのメッセージを、1つのコンシューマだけが処理する
                    → 仕事の割り振り(ジョブキュー)

トピック(1対多)    1つのメッセージを、購読している全員が受け取る
                    → 出来事の通知(Pub/Sub)

Google Cloud なら Pub/Sub(トピック + サブスクリプション)、 順序付きの大量ストリームなら Kafka、 単純なタスク実行なら Cloud Tasks といった選択肢があります。

仕組みは違っても、以下の設計上の論点はすべてに共通します。

配信保証 — ここが最重要

保証意味
at-most-once最大1回。失われることがある
at-least-once最低1回。重複することがある ← 実務のほとんどはこれ
exactly-onceちょうど1回。分散システムでは、条件付きでしか成立しない
重複は必ず起きる前提で作る

なぜ重複するのか。

1. コンシューマがメッセージを受け取る
2. 処理を完了する(DB を更新した)
3. 「処理しました」と ack を返す ← ここで落ちる
4. ack が返らないので、キューは「未処理」と判断して再配信する
5. 同じメッセージを2回処理してしまう

2 と 3 の間で落ちる可能性がある以上、重複は原理的に避けられません。

だから、冪等性が必須になります(マイクロサービスの実務)。

// メッセージ ID で処理済みを記録し、2回目は何もしない
func (h *Handler) Handle(ctx context.Context, m Message) error {
    ok, err := h.store.MarkProcessed(ctx, m.ID)   // 一意制約で重複を弾く
    if err != nil { return err }
    if !ok { return nil }                          // 既に処理済み。正常終了する
    return h.process(ctx, m)
}

「メールが2通届く」「請求が2回発生する」の原因は、ほぼこれです。

順序

多くのメッセージング基盤は、順序を保証しません。

送信: [作成] [更新] [削除]
受信: [更新] [削除] [作成]   ← ありうる

対策は2つです。

1. 順序キーを使う(同じキーのメッセージは順番に配信される)
   → ただしスループットは落ちる(並列に処理できない)

2. 順序に依存しない設計にする  ← こちらが望ましい
   → メッセージにバージョンやタイムスタンプを持たせ、古いものは無視する
// 順序が入れ替わっても壊れない
if incoming.Version <= current.Version {
    return nil   // 古い更新なので無視する
}

失敗したメッセージをどうするか

処理に失敗
   ↓ 再配信(バックオフつき。性能と負荷対策)
再び失敗
   ↓ 規定回数を超えたら
[デッドレターキュー(DLQ)] へ移す
DLQ を作って、誰も見ていない

これは非常によくあります。

DLQ にメッセージが溜まっていることは、 処理されるべき仕事が失われているという意味です。

□ DLQ の件数をメトリクスに出し、1件でも入ったらアラートを出す
□ DLQ のメッセージを調査・修正して再投入する手順を決めておく

DLQ は「捨て場」ではなく「要調査の箱」です。

毒メッセージ(poison message)

処理すると必ず失敗するメッセージが、リトライを繰り返して キュー全体を詰まらせることがあります。

順序保証付きのキューでは、1件が詰まると後続がすべて止まります。

□ リトライ回数の上限を必ず設定する
□ 上限に達したら DLQ へ移し、後続を進める
□ パースできないメッセージは、リトライせず即 DLQ へ

DB とメッセージの整合 — Outbox パターン

新人がまず気づかない、しかし本番で必ず問題になる論点です。

// 一見、正しく見える
tx.Save(order)          // DB に保存
tx.Commit()
publisher.Publish(msg)  // ← ここで落ちたら? メッセージが送られない

逆の順序にすると、今度は 「メッセージは送ったが DB のコミットに失敗した」が起きます。

DB とメッセージブローカーという2つのシステムに、 アトミックに書き込むことはできません。

解決策が Outbox パターンです。

1. 業務データとメッセージを、同じトランザクションで DB に書く
   (outbox テーブルに1行入れるだけ)
2. 別のプロセスが outbox を読み、メッセージを送信する
3. 送信できたら outbox の行に送信済みを記録する
BEGIN;
INSERT INTO orders (...) VALUES (...);
INSERT INTO outbox (id, topic, payload) VALUES (...);   -- 同じトランザクション
COMMIT;

これで、「DB は更新されたのに通知されていない」が起きなくなります。 代わりに送信が重複しうるので、受け取る側の冪等性が前提になります。

コンシューマの設計

□ 処理時間の上限を意識する(ack の期限を超えると再配信される)
□ 長い処理は、途中で ack の延長をするか、さらに分割する
□ 並列度を決める(増やせば速いが、DB の接続数を食う。性能と負荷対策)
□ 1メッセージの処理は、できるだけ小さく
□ メッセージには ID だけを入れ、実体は DB から取る
メッセージに大きなデータを入れない
悪い: メッセージに、注文の全データと商品リストを入れる
良い: メッセージには orderID だけを入れ、受け取った側が DB から読む

理由は3つあります。

  • サイズ上限がある
  • メッセージの中身が古くなる(送信後に注文が変更されたら?)
  • スキーマ変更のたびに、送信側と受信側の両方を直す必要が出る

ただし、「その時点の値」が必要な場合は例外です (監査ログや、後から変わってはいけない金額など)。

メッセージのスキーマ

スキーマと RPCの後方互換性が、そのまま当てはまります。

□ フィールドの追加は安全、削除と型変更は危険
□ 送信側と受信側は、同時にはデプロイされない
□ 未知のフィールドを受け取っても壊れないようにする
□ protobuf でスキーマを定義するのが確実

キューには「まだ処理されていない古い形式のメッセージ」が残っています。 デプロイ直後に、旧形式のメッセージが新しいコードに届くことを想定してください。

監視

□ 未処理メッセージの件数(滞留)
□ 最も古い未処理メッセージの経過時間  ← 最重要
□ DLQ の件数
□ 処理のエラー率と処理時間
「静かに遅れる」のが非同期の怖さ

同期呼び出しなら、遅ければ利用者がすぐ気づきます。

非同期は違います。処理が止まっていても、システムは正常に見えます。 気づくのは「メールが来ない」という問い合わせが来た時です。

だから滞留の監視は必須です。 「未処理の最古メッセージが10分を超えたらアラート」のように、 時間で監視するのが実務的です(監視とオンコール)。

注文確定時に確認メールを送る処理を、キュー経由の非同期にしました。テスト中、まれにメールが2通届きます。どう対応しますか。

実務の落とし穴まとめ

  1. 冪等でないコンシューマ — 重複は必ず起きる。二重請求・二重送信の原因
  2. 順序を前提にした設計 — 保証されない。バージョンで判定する
  3. DLQ を誰も見ていない — 仕事が静かに失われている
  4. リトライ上限が無い — 毒メッセージがキューを詰まらせる
  5. DB 更新とメッセージ送信を別々に行う — Outbox パターンを使う
  6. メッセージに大きなデータを入れる — 古くなる、サイズ制限、変更が波及
  7. 滞留を監視していない — 止まっていても正常に見える
  8. 旧形式のメッセージを考慮しない — デプロイ直後に壊れる
  9. すぐ結果が必要な処理を非同期にする — 複雑になるだけ

まとめ

  • 非同期は独立性を買う手段。代償として複雑さと「いつ終わるか分からない」を受け入れる
  • 実務のほとんどは at-least-once。重複は仕様であり、冪等性が必須
  • 順序は保証されない。順序キーを使うか、順序に依存しない設計にする
  • 失敗はバックオフ付きリトライ → 上限 → DLQ。DLQ は要調査の箱
  • DB とブローカーにアトミックには書けない。Outbox パターンで解決する
  • メッセージにはID を入れ、実体は DB から取る(時点の値が必要な場合を除く)
  • スキーマは後方互換を守る。キューには旧形式が残っている
  • 滞留を時間で監視する。非同期の障害は静かに進行する

公式ドキュメント

迷ったら一次情報に戻ってください。

章末問題

注文をDBに保存した後、Pub/Sub にイベントを publish するコードがあります。まれに「注文はあるのに在庫が引き当てられていない」状態が発生します。原因と対策は?

次の章では、このキューを実際に動かしている製品—— RabbitMQ と Kafka を、それぞれ何に向くかまで見ていきます。

読み終わったら記録しておくと、目次で進み具合が分かります。