メッセージキュー
サービス同士を「あとで確実に処理する」で繋ぐキューを作る。メッセージは追記ログに積まれ、消費者が読んだ位置をオフセットとして持つ。配送保証は処理とオフセット確定のどちらを先にやるかで決まり、確定が先なら取りこぼし(at-most-once)、処理が先なら重複(at-least-once)になる。実質 1 回は at-least-once と冪等な消費者で作り、クラッシュを起こしながら確かめる。
この章で作るもの
同期呼び出し(RPC)は相手が生きていないと成立しない。メッセージキューは「頼みごとをログに置いておき、相手が後で取りに来る」非同期の繋ぎ方。送る側と処理する側を切り離し(疎結合)、負荷の波を吸収し、相手が落ちても取りこぼさない。題材は Kafka のようなログ型キュー。
順に見ていく。
- ログ + オフセット: メッセージは WAL と同じ追記専用ログに積まれる。ブローカはメッセージを消さない。消費者が「どこまで読んだか」のオフセットを進めるだけ。だから複数の消費者が独立に、好きな位置から読める
- 配送保証: メッセージを「処理」し、オフセットを「確定」する。この順序で保証が決まる。確定が先なら取りこぼし(at-most-once)、処理が先なら重複(at-least-once)
- 実質1回: at-least-once + 冪等な消費者(同じ処理を2回やっても結果が変わらない)。再配送されても副作用は1回きりになる
ログ + オフセット
キューといっても、中身は追記ログ1本と、消費者ごとのオフセット1つ。ブローカはメッセージを消さず、消費者が読んだ位置を覚えているだけ。この単純さが強い。複数の消費者が同じログを別々の速度で読めるし、消費者が落ちても、記録したオフセットから読み直せる。
ブローカのログ(追記専用・消えない)
┌───┬───┬───┬───┬───┬───┐
│ 0 │ 1 │ 2 │ 3 │ 4 │ 5 │← 末尾に Publish
└───┴───┴───┴───┴───┴───┘
▲ ▲
消費者B(2) 消費者A(4) ← それぞれのオフセット配送保証: 取りこぼしか、重複か
分かれ目は「メッセージの処理と、オフセットの確定の、どちらを先にやるか」になる。ここでクラッシュが起きると差が出る:
// Poll は未読を最大 max 件、配送保証に従って処理する。処理した件数を返す。
//
// AtMostOnce : オフセットを先に進めてから処理する(処理前に落ちると取りこぼす)
// AtLeastOnce: 処理してからオフセットを進める(確定前に落ちると再配送される)
func (c *Consumer) Poll(max int) int {
batch := c.broker.fetch(c.committed, max)
if c.semantics == AtMostOnce {
c.committed += len(batch) // 先に確定
for _, m := range batch {
c.handle(m) // ここで落ちると、この batch は再配送されない=消失
}
return len(batch)
}
// AtLeastOnce
for _, m := range batch {
c.handle(m) // 先に処理
}
c.committed += len(batch) // 後で確定。ここに来る前に落ちると再配送=重複
return len(batch)
}- at-most-once(確定 → 処理): 先にオフセットを進める。処理の前に落ちると、そのメッセージは確定済み扱いで再配送されず、失われる。重複はしないが取りこぼす
- at-least-once(処理 → 確定): 先に処理する。確定の前に落ちると、次回再配送され、同じメッセージを2回処理しうる。取りこぼさないが重複する
at-most-once : [確定]──💥落ちる──[処理] → 処理されず消失(取りこぼし)
at-least-once: [処理]──💥落ちる──[確定] → 確定が飛ぶ → 次回 再配送(重複)動かす
下のデモで「発行」→「消費」でオフセットが進む様子を見て、「クラッシュ消費」(確定の前後で中断)を押す。at-most-once だと取りこぼし、at-least-once だと次の「消費」で再配送され、同じメッセージが2回処理される(下段の効果ログに重複が出る)。
発行して消費してみる。『クラッシュ消費』で確定前に落ちると挙動が分かれる
実質1回: 冪等な消費者
「取りこぼさない」を選ぶと at-least-once、つまり重複は避けられない。ならば重複しても平気な消費者にすればいい。同じメッセージを2回処理しても結果が変わらない性質を**冪等(idempotent)**という。
やり方はシンプル。メッセージに一意なキー(注文ID・イベントID)を持たせ、消費者は「このキーはもう処理したか」を覚えておき、既処理なら捨てる:
// Handle は1件を受ける。初見なら副作用を記録し、既見(重複)なら捨てる。
// Consumer の handle にそのまま渡せる。
func (s *IdempotentSink) Handle(m Message) {
if _, dup := s.seen[m.Key]; dup {
s.Duplicates++
return // すでに処理済みのキー。再配送なので副作用は起こさない
}
s.seen[m.Key] = struct{}{}
s.Delivered = append(s.Delivered, m)
}デモで**「冪等 ON」**にしてから同じクラッシュ→再配送を起こすと、2回目の配送は「×重複」として捨てられ、実質適用は各メッセージ1回きりに保たれる。これが at-least-once + 冪等 = 実質1回(effectively-once)。分散システムで「取りこぼさない・重複させない」を両立する、最も実用的な落としどころ。
なぜ「厳密に1回」を狙わないのか
送信と確定は別々の操作で、その間はいつでも落ちうる。「厳密に1回だけ届く」を通信レベルで完全保証するのは一般に不可能(2将軍問題)。だから実務は「少なくとも1回届く」を保証し、受け手を冪等にして重複を無害化する。Kafka の "exactly-once" も、内部はこの冪等 + トランザクションで実現している。
設計の観点: キュー設計の見積もり
「注文イベントを 5万/秒でさばく。どう設計する?」。キューの問いは throughput と順序で答える:
- パーティション数: 1本のログ(パーティション)は順序を保てるが、1消費者しか並列に処理できない。スループットを上げるにはキーでパーティション分割し、複数消費者で並列化する。5万/秒 ÷ 1消費者1万/秒 = 5パーティション以上
- 順序保証の範囲: 全体の順序は諦め、「同じ注文IDのイベントは同じパーティション」に寄せてキー単位の順序だけ守る。パーティション割り当てはまさに コンシステントハッシュ
- 冪等性: at-least-once 前提で、消費側を必ず冪等に。「二重課金しない」は再送耐性で担保する
- バックプレッシャと保持: 消費が詰まったらログが伸びる。保持期限(retention)とディスク量、遅延(lag)を監視項目に
メリット・デメリットと実例
| 型 | 保持 | 順序 | 主な用途 | 実例 |
|---|---|---|---|---|
| ログ型 | 期限まで残る(再読可) | パーティション内 | イベント基盤・ストリーム処理 | Apache Kafka、Redis Streams、AWS Kinesis |
| キュー型 | 消費で消える | 基本なし(FIFO は別) | ジョブ分配・タスクキュー | RabbitMQ、AWS SQS、NATS |
裏どり:
- Apache Kafka: 追記ログをパーティションに分け、消費者はオフセットをコミットして読む。この章の
Broker/Consumerの直接の元。retention 期間内なら何度でも読み直せる - AWS SQS: キュー型。受信したメッセージは可視性タイムアウトの間だけ隠れ、削除しなければ再配送(at-least-once)。標準キューは順序保証なし、FIFO キューは別モード
- RabbitMQ: ブローカがルーティングし、ack しないと再配送。プッシュ型でジョブ分配向き
- いずれも既定は at-least-once。「厳密に1回」を謳うものも、実体は冪等/トランザクションによる重複無害化
簡略化したこと
- 1パーティション・1ログのみ。パーティション分割/コンシューマグループは無し(キーで分割し コンシステントハッシュ でノードへ割る、が実物の形)
- 永続化・ネットワーク・保持期限は無し(ログはメモリ、消えない)
- 冪等は「Key を覚える」だけ。実務は処理結果側の UNIQUE 制約や処理済み表で担保する
- 順序保証は1ログ内のみ。複数パーティション間の全順序は扱わない
参考資料
- Kleppmann, Designing Data-Intensive Applications 11章(Stream Processing)
- Kafka ドキュメント: Consumer offset / delivery semantics / exactly-once
- Amazon SQS ドキュメント: 可視性タイムアウトと at-least-once
- 実装: messaging/queue