【実務・中級編】 Stream読み取りコマンド – Redis

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アラートからあなたを解放してくれる。コードレビューでこれらが抜けていたら、容赦なく差し戻してほしい。

コメント

タイトルとURLをコピーしました