Redis Stream徹底解剖:XREAD / XREADGROUPの極限と実務設計のアンチパターン
テックリードの私だ。コードレビューや設計レビューで、未だに「とりあえず `XREAD` でポーリングすればいいか」「コンシューマーグループは難しそうだから単一のストリームをみんなでゴリゴリ読もう」といった、甘えた設計を見かける。
RedisのStreamは、単なる「ちょっと高機能なリスト」ではない。Kafkaの思想を洗練させ、Redisのインメモリの猛威と組み合わせた本気のエントリ・メッセージング・データ構造だ。
今回は、`XREAD` と `XREADGROUP` の挙動の裏側、そして実務の現場でシステムを沈めないための堅牢な設計パターンを、容赦なくコードとアーキテクチャの観点から解説する。
—
1. 根本思想の理解:なぜPub/SubやListではなくStreamなのか?
まず前提を叩き込んでおく。
- Pub/Sub: 射撃して終わり(Fire-and-forget)。オフラインのクライアントはデータを取りこぼす。永続化されない。
- List (LPRP/RPOP): 永続化はされるが、複数コンシューマーで「メッセージの奪い合い」になり、どのコンシューマーが処理に失敗したかの追跡(Ack)が自前実装になる。
- Stream: オフセット管理、ACK(確認応答)、コンシューマーグループによるロードバランス、未処理メッセージ(Pending Entries List)の再分配をRedisコアが面倒を見てくれる。
これを自前でRDBやRedisのHash/Listで実装しようとした瞬間から、あなたのシステムは死へのカウントダウンを始める。Redis Streamを使え。
—
2. `XREAD`:単一ストリーム / マルチストリームのストリーム・リーダー
`XREAD` は、指定したIDよりも新しいメッセージを読み出すためのコマンドだ。ブロック(ブロッキング待機)機能を備えており、Pub/Subのようなリアルタイムプッシュ受信を模倣できる。
基本構文と罠
ストリーム “sensor-log” から、ID “0-0” より大きいメッセージをブロックなしで最大2件取得
XREAD COUNT 2 STREAMS sensor-log 0-0
実務における `XREAD` の限界
`XREAD` は強力だが、「1つのメッセージを複数のコンシューマーで分散処理する(ロードバランス)」ことはできない。
`XREAD` を使った複数ワーカー構成にすると、すべてのワーカーが「全く同じメッセージを重複して受信」することになる。これは分散キューではなく、ただの全配信(Fan-out)だ。
したがって、`XREAD` を実務で使うべきユースケースは厳格に限定される:
1. 単一のプロセスがすべてのイベントを監視・処理する(イベントロガー、単一のダッシュボード等)。
2. 「ファンアウト(各サービスが全てのイベントのコピーを受け取る必要がある)」を意図的に行いたい場合。
—
3. `XREADGROUP`:真の分散メッセージングとコンシューマーグループ
複数台のワーカー(コンシューマー)で巨大なメッセージストリームを分担処理したいのであれば、選択肢は `XREADGROUP` 一択だ。
アーキテクチャの核心:PEL(Pending Entries List)
コンシューマーグループを使うと、Redisサーバー側で以下の状態が管理される。
1. Last Delivered ID: グループ内で「誰かに割り当てられた」最後のメッセージID。
2. PEL (Pending Entries List): 「特定のコンシューマーに配信されたが、まだ `XACK`(処理完了確認)されていないメッセージ」のリスト。
このPELがあるおかげで、ワーカーが途中でクラッシュしてもメッセージがロストしないという堅牢性が担保される。
実践:グループの作成から読み取りまでのライフサイクル
1. コンシューマーグループの作成 (存在しない場合はストリームも同時に作成)
“my-group” を作成。 “$” は「これ以降の新着メッセージのみ対象」を意味する。「0」にすると過去の未読分も対象になる。
XGROUP CREATE sensor-log my-group $ MKSTREAM
2. コンシューマー “worker-1” としてメッセージを読み取る
BLOCK 5000: メッセージがない場合、最大5秒間ブロックする
COUNT 10: 一度に最大10件取得
> : 「このグループにまだ割り当てられていない、新しいメッセージ」を要求する特殊なID
XREADGROUP GROUP my-group worker-1 COUNT 10 BLOCK 5000 STREAMS sensor-log >
3. メッセージの処理完了を通知 (XACK)
これにより、PELから該当メッセージが削除される
XACK sensor-log my-group 1672531200000-0
—
4. チーフアーキテクトが伝授する「実務で踏み抜く地雷」と設計パターン
ここからが本題だ。綺麗事のドキュメントには書いていない、現場でシステムを崩壊させる罠と、その回避策を授ける。
地雷1: 「`>`」指定の罠とリカバリの欠如
`XREADGROUP` で `>` を指定すると、「まだどのコンシューマーにも渡っていない新規メッセージ」が取れる。しかし、ワーカーがメッセージを取得した瞬間にクラッシュしたらどうなるか?
メッセージはPELに残り続け、永遠に `>` では取得できない。
【解決策】定期的なPELの棚卸し (`XPENDING` と `XCLAIM`)
システムには必ず「ゾンビ・ワーカー(死んだのに処理途中のままのワーカー)」を救出する掃除屋(Reclaimer)プロセスを常駐させろ。
1. PELに溜まっている古い未処理メッセージを検知する
引数: ストリーム名、グループ名、最小アイドル時間(ミリ秒)、取得上限
XPENDING sensor-log my-group – + 10 worker-1
2. 死んだワーカーのメッセージを、生きているワーカー(例: worker-2)に強制移譲する
60000ミリ秒(1分)以上放置されているメッセージを worker-2 のものにする
XCLAIM sensor-log my-group worker-2 60000 1672531200000-0
この `XPENDING` -> `XCLAIM` のループをバックグラウンドで回すことが、ロバストなメッセージング基盤の絶対条件だ。
—
地雷2: メッセージ肥大化によるメモリ爆発
Redisはインメモリデータベースだ。RDBやKafkaのように、数テラバイトのデータをディスクに退避させながら無限に保持できるわけではない。
ストリームにデータを入れ続け、誰も `XDEL` しなかったり、コンシューマーグループのポインタが進まないと、メモリは確実に枯渇する。
【解決策】`MAXLEN` による自動トリミング
メッセージの追加時 (`XADD`) には、必ず `MAXLEN`(または近似的トリミングの `~`)を指定せよ。
ストリームの長さを最大10,000件に維持し、古いものは自動削除する
XADD sensor-log MAXLEN ~ 10000 temperature 23.5 humidity 60
※注意:コンシューマーグループが遅れている状態で `MAXLEN` で古いデータを無理やり消すと、「まだ誰も読んでいないのに消える」という悲劇が起きる。グループの処理速度と `MAXLEN` のサイズ設計は、インフラエンジニアとしての腕の見せ所だ。
—
5. まとめ:プロダクション品質のコードを書くために
Redis Streamの `XREAD` / `XREADGROUP` は、正しく扱えばKafka顔負けの軽量で高速なメッセージングシステムを構築できる。
設計時のチェックリストを心に刻んでおけ:
1. 単なるブロッキング購読なら `XREAD`、分散・耐久性が必要なら必ず `XREADGROUP` を使っているか?
2. `XREADGROUP` で取得したメッセージに対し、処理成功時に必ず `XACK` を呼んでいるか?
3. クラッシュしたワーカーの放置メッセージを救出する `XCLAIM` の仕組み は実装されているか?
4. メモリ枯渇を防ぐための `MAXLEN` によるトリミング は設計されているか?
これらを満たしたコードだけが、深夜のPagerDutyアラートからあなたを解放してくれる。コードレビューでこれらが抜けていたら、容赦なく差し戻してほしい。
コメント