Skip to content

Pub/Sub

「1 つのメッセージを購読者全員に配る」ファンアウトを作る。発行は 1 回で、購読者はそれぞれ全メッセージを受け取る(1 件を 1 人が奪うキューの逆だ)。購読者ごとに読んだ位置(カーソル)を独立に持つので、遅い購読者がいても他は止まらない。購読を始める位置も、過去から再生するか以降だけかを選べる。キューとの配り方の違いを、購読者を足しながらデモで見る。

この章で作るもの

メッセージキューは「1件を1人が取る」点対点だった。ジョブを複数ワーカで分担するのに向く。Pub/Sub はその逆で「1件を全員に配る」ファンアウト。1つのイベント(注文された、が更新された)を、在庫・通知・分析…と複数の関心先に一斉配信するのに向く。

順に見ていく。

  1. ファンアウト: 発行は1回。トピックの購読者はそれぞれ独立に全メッセージを読む。キューが「奪い合い」なら、Pub/Sub は「配り」
  2. 独立カーソル: 購読者ごとに「どこまで読んだか」を持つ。だから遅い購読者がいても、他の購読者は影響を受けず自分のペースで進める
  3. 購読開始位置: 新しく購読を始めたとき、過去のメッセージを再生するか(FromBeginning)、以降だけ受け取るか(FromNow)

キューとの違い: 奪い合いか、配りか

同じ「ログ + オフセット」でも、配り方が逆になる。キューは1件を1人が取って消費を分担する。Pub/Sub は1件を購読者全員が受け取る。

  キュー(点対点)                    Pub/Sub(ファンアウト)
   [m1][m2][m3]                       [m1][m2][m3]
     │   │   │                        ├──全部──▶ 購読者A
     ▼   ▼   ▼                        ├──全部──▶ 購読者B
   W1  W2  W3 (1件を1人)              └──全部──▶ 購読者C (各自が全部)
点対点(キュー)とファンアウト(Pub/Sub)。キューは1メッセージを1ワーカが取り、負荷を分担する。Pub/Sub は1メッセージを全購読者へ配る。同じログでも、カーソルを『共有して奪い合う』か『各自持って全部読む』かの違い

ファンアウトと独立カーソル

購読者はそれぞれ自分のカーソルを持ち、そこから未読を読む。発行は1回でも、購読者の数だけ「読む主体」がいる:

go
// Poll は未読を最大 max 件返し、カーソルを進める。購読者ごとに独立なので、
// ある購読者が読まなくても(遅くても)、他の購読者には影響しない。
func (s *Subscription) Poll(max int) []Message {
	log := s.broker.topics[s.topic]
	end := s.cursor + max
	if end > len(log) {
		end = len(log)
	}
	if s.cursor >= end {
		return nil
	}
	out := make([]Message, end-s.cursor)
	copy(out, log[s.cursor:end])
	s.cursor = end
	return out
}

カーソルが購読者ごとに独立なので、遅い購読者は他に影響しない。分析用の重いバッチが遅れていても、通知用の軽い購読者はリアルタイムに進める。ただし遅い購読者のバックログは伸び続ける。ここが後述のバックプレッシャの論点になる。

動かす

下のデモで**「発行」を押すと、購読者全員のバックログが同時に +1 される(ファンアウト)。各購読者の「受信」**は独立で、1人が受信しても他のカーソルは動かない。片方だけ受信して、もう片方のバックログが残るのを確かめてほしい。

デモPub/Sub(トピックのファンアウト)購読者 2 / ログ 0件
トピック "news" のログ(追記され、消えない)
(空)
購読者 1FromBeginning未読 0
(受信なし)
購読者 2FromBeginning未読 0
(受信なし)

発行すると購読者全員のバックログが増える(ファンアウト)。各自が自分のペースで受信する

発行1回 → 購読者全員のバックログ +1(ファンアウト)。キューは1件を1人が奪う各購読者は独立カーソル。遅い購読者がいても他は止まらない

購読開始位置: 過去を見るか

新しい購読者が現れたとき、過去のメッセージを届けるかが設計の分かれ目:

go
// Subscription は1つのトピックを、自分のカーソルで読む購読。
type Subscription struct {
	broker *Broker
	topic  string
	cursor int // 次に読む offset
}

// Subscribe はトピックの購読を作る。FromNow なら現在の末尾から始める。
func (b *Broker) Subscribe(topic string, start StartAt) *Subscription {
	cursor := 0
	if start == FromNow {
		cursor = len(b.topics[topic]) // 今ある分は飛ばす
	}
	return &Subscription{broker: b, topic: topic, cursor: cursor}
}
  • FromBeginning: 先頭から再生。「これまでの全イベントを処理し直したい」新しい分析パイプラインなど(ログが残っている=durable)
  • FromNow: 購読時点の末尾から。「今から起きることだけ知りたい」通知など(過去は不要=ephemeral)

デモで「購読者+ (FromNow)」を、いくつか発行した後に追加すると、そのカーソルが末尾から始まり過去分を飛ばすのが分かる。Kafka の「オフセットをどこから読むか(earliest / latest)」がこれ。

設計の観点: ファンアウトの使いどころ

「注文イベントを在庫・通知・分析の3系統に流したい」。Pub/Sub の出番:

  • キューか Pub/Sub か: 「1つの仕事を分担」ならキュー、「1つのイベントを複数の関心先へ」なら Pub/Sub。注文確定を在庫と通知と分析がそれぞれ処理するなら後者
  • 遅い購読者(slow consumer): 独立カーソルの代償で、遅い購読者のバックログは伸びる。放置するとログが膨らむ。対策は保持期限で古いものを捨てる/購読者を落とす/バックプレッシャで発行を絞る
  • 配送保証は別問題: ファンアウトしても各購読者への配送は at-least-once が普通。受け手はやはり冪等に(メッセージキューの冪等)
  • トピック設計: 粗すぎると要らないメッセージまで届き、細かすぎると購読管理が複雑。イベントの種類でトピックを切る

メリット・デメリットと実例

過去の再生主な用途実例
ログ型(durable)できる(FromBeginning)イベントソーシング・再処理Kafka(コンシューマグループ)、Redis Streams
ブロードキャスト型(ephemeral)できない(FromNow 相当)リアルタイム通知・チャットRedis Pub/Sub、MQTT、WebSocket ファンアウト
マネージド設定次第疎結合な通知配信AWS SNS、Google Cloud Pub/Sub

裏どり:

  • Kafka: 1つのログを、コンシューマグループ単位でファンアウトする。同じグループ内は点対点(分担)、グループを分ければ各グループが全メッセージを受け取る(Pub/Sub)。この章の「独立カーソル」はグループごとのオフセットに当たる
  • Redis Pub/Sub: ephemeral。購読していない間のメッセージは届かない(FromNow のみ)。Redis Streams は durable で過去も読める
  • MQTT: IoT 定番。トピック階層とワイルドカード購読(sensor/+/temp)を持つ
  • AWS SNS: トピックに publish すると、購読している SQS キュー・Lambda・HTTP へ一斉配信

簡略化したこと

  • ワイルドカード購読(MQTT の階層トピック)は無し。トピックは完全一致
  • バックプレッシャ・保持期限・遅い購読者の drop ポリシーは無し(ログは伸び続ける)
  • ネットワーク・永続化は無し(全てメモリ)。配送保証(冪等・再送)は メッセージキュー 側の話
  • 購読解除やトピック削除は最小限(デモのみ)

参考資料

  • Kleppmann, Designing Data-Intensive Applications 11章(Stream Processing)
  • Kafka: コンシューマグループとオフセット
  • AWS SNS / MQTT 仕様(トピックとファンアウト)
  • 実装: messaging/pubsub