非同期処理とメッセージング
この部の 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 の件数をメトリクスに出し、1件でも入ったらアラートを出す
□ DLQ のメッセージを調査・修正して再投入する手順を決めておく
DLQ は「捨て場」ではなく「要調査の箱」です。
処理すると必ず失敗するメッセージが、リトライを繰り返して キュー全体を詰まらせることがあります。
順序保証付きのキューでは、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通届きます。どう対応しますか。
実務の落とし穴まとめ
- 冪等でないコンシューマ — 重複は必ず起きる。二重請求・二重送信の原因
- 順序を前提にした設計 — 保証されない。バージョンで判定する
- DLQ を誰も見ていない — 仕事が静かに失われている
- リトライ上限が無い — 毒メッセージがキューを詰まらせる
- DB 更新とメッセージ送信を別々に行う — Outbox パターンを使う
- メッセージに大きなデータを入れる — 古くなる、サイズ制限、変更が波及
- 滞留を監視していない — 止まっていても正常に見える
- 旧形式のメッセージを考慮しない — デプロイ直後に壊れる
- すぐ結果が必要な処理を非同期にする — 複雑になるだけ
まとめ
- 非同期は独立性を買う手段。代償として複雑さと「いつ終わるか分からない」を受け入れる
- 実務のほとんどは at-least-once。重複は仕様であり、冪等性が必須
- 順序は保証されない。順序キーを使うか、順序に依存しない設計にする
- 失敗はバックオフ付きリトライ → 上限 → DLQ。DLQ は要調査の箱
- DB とブローカーにアトミックには書けない。Outbox パターンで解決する
- メッセージにはID を入れ、実体は DB から取る(時点の値が必要な場合を除く)
- スキーマは後方互換を守る。キューには旧形式が残っている
- 滞留を時間で監視する。非同期の障害は静かに進行する
公式ドキュメント
迷ったら一次情報に戻ってください。
| 対象 | リンク |
|---|---|
| Google Cloud Pub/Sub(日本語) | https://cloud.google.com/pubsub/docs?hl=ja |
| Amazon SQS(日本語) | https://docs.aws.amazon.com/ja_jp/AWSSimpleQueueService/latest/SQSDeveloperGuide/ |
| Enterprise Integration Patterns | https://www.enterpriseintegrationpatterns.com/ |
| microservices.io: Transactional Outbox | https://microservices.io/patterns/data/transactional-outbox.html |
章末問題
注文をDBに保存した後、Pub/Sub にイベントを publish するコードがあります。まれに「注文はあるのに在庫が引き当てられていない」状態が発生します。原因と対策は?
次の章では、このキューを実際に動かしている製品—— RabbitMQ と Kafka を、それぞれ何に向くかまで見ていきます。