Skip to content

メッセージキュー

サービス同士を「あとで確実に処理する」で繋ぐキューを作る。メッセージは追記ログに積まれ、消費者が読んだ位置をオフセットとして持つ。配送保証は処理とオフセット確定のどちらを先にやるかで決まり、確定が先なら取りこぼし(at-most-once)、処理が先なら重複(at-least-once)になる。実質 1 回は at-least-once と冪等な消費者で作り、クラッシュを起こしながら確かめる。

この章で作るもの

同期呼び出し(RPC)は相手が生きていないと成立しない。メッセージキューは「頼みごとをログに置いておき、相手が後で取りに来る」非同期の繋ぎ方。送る側と処理する側を切り離し(疎結合)、負荷の波を吸収し、相手が落ちても取りこぼさない。題材は Kafka のようなログ型キュー。

順に見ていく。

  1. ログ + オフセット: メッセージは WAL と同じ追記専用ログに積まれる。ブローカはメッセージを消さない。消費者が「どこまで読んだか」のオフセットを進めるだけ。だから複数の消費者が独立に、好きな位置から読める
  2. 配送保証: メッセージを「処理」し、オフセットを「確定」する。この順序で保証が決まる。確定が先なら取りこぼし(at-most-once)、処理が先なら重複(at-least-once)
  3. 実質1回: at-least-once + 冪等な消費者(同じ処理を2回やっても結果が変わらない)。再配送されても副作用は1回きりになる

ログ + オフセット

キューといっても、中身は追記ログ1本と、消費者ごとのオフセット1つ。ブローカはメッセージを消さず、消費者が読んだ位置を覚えているだけ。この単純さが強い。複数の消費者が同じログを別々の速度で読めるし、消費者が落ちても、記録したオフセットから読み直せる。

  ブローカのログ(追記専用・消えない)
  ┌───┬───┬───┬───┬───┬───┐
  │ 0 │ 1 │ 2 │ 3 │ 4 │ 5 │← 末尾に Publish
  └───┴───┴───┴───┴───┴───┘
            ▲           ▲
        消費者B(2)   消費者A(4)     ← それぞれのオフセット
ログ + オフセット。メッセージは末尾に追記される。消費者Aはオフセット4、消費者Bはオフセット2まで読んだ。ブローカは消さないので、両者が独立に進める。落ちても自分のオフセットから再開できる

配送保証: 取りこぼしか、重複か

分かれ目は「メッセージの処理と、オフセットの確定の、どちらを先にやるか」になる。ここでクラッシュが起きると差が出る:

go
// 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 は処理が先なので、確定前に落ちると再配送され重複する。『取りこぼさない』と『重複しない』はクラッシュ点の前後で交換になる

動かす

下のデモで「発行」→「消費」でオフセットが進む様子を見て、「クラッシュ消費」(確定の前後で中断)を押す。at-most-once だと取りこぼし、at-least-once だと次の「消費」で再配送され、同じメッセージが2回処理される(下段の効果ログに重複が出る)。

デモメッセージキュー(ログ + オフセット)確定 0/0
at-most-onceat-least-once
冪等 OFF
ブローカのログ(追記され、消えない)
(空)
消費者オフセット: 0(ここから先が未読)
処理の効果 — handle 呼び出し 0回 / 実質適用 0件
(まだ処理なし)

発行して消費してみる。『クラッシュ消費』で確定前に落ちると挙動が分かれる

確定済み(消費者が読み終えた)at-least-once + 冪等で「取りこぼさない・重複させない」= 実質1回

実質1回: 冪等な消費者

「取りこぼさない」を選ぶと at-least-once、つまり重複は避けられない。ならば重複しても平気な消費者にすればいい。同じメッセージを2回処理しても結果が変わらない性質を**冪等(idempotent)**という。

やり方はシンプル。メッセージに一意なキー(注文ID・イベントID)を持たせ、消費者は「このキーはもう処理したか」を覚えておき、既処理なら捨てる:

go
// 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