Redis Streamコンシューマーグループ:本番障害を防ぐためのアーキテクチャと実装パターン
こんにちは。テクニカルリードの私だ。
今日のコードレビューで、また「とりあえず`XREAD`でポーリングしてメッセージ処理しました。複数ワーカーで動かすために適当に分散してます」という実装を見た。
……待て。その設計、本当にプロダクションの負荷に耐えられるのか? ワーカーが突然死したらメッセージはロストしないのか? 処理途中のメッセージが宙に浮いたまま放置されていないか?
KafkaやRabbitMQのノリでRedisのStreamを雑に扱うと、必ず痛い目を見る。Redis Streamは強力だが、その真価は「コンシューマーグループ(Consumer Groups)」を正しく理解し、適切なAck/Pending管理を組み込んだ瞬間に初めて発揮される。
今回は、RedisのStream型におけるコンシューマーグループのメカニズムを解剖し、実務で絶対に破綻しない堅牢なメッセージングシステムの構築方法を叩き込む。
—
1. なぜ「Pub/Sub」でも「ListのBLPOP」でもなく「Stream + コンシューマーグループ」なのか?
まず前提を合わせよう。
- Pub/Sub: 揮発性。メッセージの永続化なし。購読していない瞬間に飛んだメッセージは永遠に消える。
- List (LPRP/RPOP等): 永続化はできるが、単一のキューを複数ワーカーで奪い合うため、「どのワーカーが処理中で、どれが未処理か」の状態管理(Ack)が自前実装になる。クラッシュ時のリカバリが地獄。
- Stream + コンシューマーグループ: Kafkaにインスパイアされた強力な機能。メッセージの永続化、スケーラブルな負荷分散(パーティショニング)、そして「誰がどのメッセージを処理していて、まだ完了していないか」のトラッキング(Pending Entries List: PEL)が完備されている。
実務で信頼性の高い非同期処理基盤やイベント駆動アーキテクチャを作るなら、選択肢はこれ一択だ。
—
2. コアコマンド群の正しい理解と実践
コンシューマーグループを操る上で、以下の4つのコマンド群の挙動を正確に把握しておく必要がある。
1. `XGROUP`: グループの作成・管理
2. `XREADGROUP`: グループ内コンシューマーによるメッセージ読み込み
3. `XACK`: メッセージ処理完了の通知(非常に重要)
4. `XPENDING`: 未完了メッセージ(PEL)の監査と管理
ステップ1: コンシューマーグループの作成 (`XGROUP CREATE`)
まずはストリームと、それを購読するグループを作る。
すでにストリームが存在しない場合は、`MKSTREAM`オプションをつけるのが定石だ。
orders ストリームに対して、processing_group という名前のグループを先頭($)から作成
$ は「これ以降に届いた新しいメッセージのみを対象とする」という意味
127.0.0.1:6379> XGROUP CREATE orders processing_group $ MKSTREAM
OK
※もし過去のメッセージも含めて全処理させたい場合は、`$` の代わりに `0` を指定する。
ステップ2: グループ内でのメッセージ読み込み (`XREADGROUP`)
ワーカーAがメッセージを読み込む。ここで重要なのは、「どのコンシューマーが読んだか」を明示することだ。
group processing_group の worker_1 として、まだどのコンシューマーにも配信されていない新しいメッセージ(>)を最大1件取得
127.0.0.1:6379> XREADGROUP GROUP processing_group worker_1 COUNT 1 STREAMS orders >
1) 1) “orders”
2) 1) “order_id”
2) “98765”
“amount”
“1500”
ここで注目すべきは `>` という特殊なIDだ。これは「このグループ内で、まだどのコンシューマーにも渡されていない未配信のメッセージ」を意味する。
このコマンドを発行した瞬間、Redisの内部ではそのメッセージがPEL(Pending Entries List:保留エントリリスト)に登録され、「worker_1が保持している状態」に変わる。
ステップ3: 処理完了の通知 (`XACK`)
ビジネスロジック(DBへの書き込みや外部API呼び出しなど)が正常に完了したら、必ず`XACK`を送る。これを行わない限り、メッセージは「処理中(未完了)」のままPELに残続け、幽霊メッセージと化す。
orders ストリームの processing_group において、該当メッセージIDの処理完了を通知
127.0.0.1:6379> XACK orders processing_group 1700000000000-0
(integer) 1
この`XACK`があって初めて、PELから該当エントリが削除され、メッセージのライフサイクルが完結する。
—
3. 実務で避けて通れない「クラッシュとリカバリ」の設計
優秀なエンジニアとそうでないエンジニアの分かれ目は、「正常系」ではなく「異常系(ワーカーの死活管理)」をどう設計しているかにある。
もし `worker_1` が `XREADGROUP` でメッセージを受け取った直後、電源断やOOM Killerで死亡したらどうなるか?
メッセージは `worker_1` のPELにぶら下がったまま、永遠に誰も処理できない状態(Zombie State)に陥る。
ここで登場するのが `XPENDING` と `XCLAIM` だ。
停滞メッセージの監査 (`XPENDING`)
グループ内で、どのコンシューマーがどれだけ未完了メッセージを抱えているかを確認する。
127.0.0.1:6379> XPENDING orders processing_group
1) (integer) 3 # 保留中の総数
2) “1699999900000-0” # 最も古い保留メッセージID
3) “1699999950000-0” # 最も新しい保留メッセージID
4) 1) 1) “worker_1”
2) “3” # worker_1 が3件持っている
詳細を見ることもできる。
まだ処理されていないメッセージのなかで、一定時間(例: 60秒)以上経過しているものを特定する
127.0.0.1:6379> XPENDING orders processing_group – + 10 worker_1
ゾンビメッセージの強奪 (`XCLAIM`)
死んだワーカー(例: `worker_1`)が抱え込んで放置しているメッセージを、生きているワーカー(例: `worker_2`)が強制的に引き継ぐ(Claimする)ためのコマンドが `XCLAIM` だ。
worker_1 が 60000ミリ秒(60秒)以上放置しているメッセージ 1700000000000-0 を、worker_2 に強制移行する
127.0.0.1:6379> XCLAIM orders processing_group worker_2 60000 1700000000000-0
【堅牢な設計パターン:Dead Letter / リカバリDAEMON】
本番環境では、メインのワーカープロセスとは別に、「リカバリ専用のバックグラウンドプロセス(またはスレッド)」を常駐させるべきだ。
このプロセスは定期的に `XPENDING` を叩き、一定時間(例: 5分)以上ACKされていないメッセージを検出し、`XCLAIM` で自グループの別ワーカーに再アサインするか、あるいはリトライ回数(独自のカウンターと組み合わせる)を超えていれば `Dead Letter Queue (DLQ)` 用の別ストリームへ退避させる。この仕組みがないシステムは、本番運用で必ずデータをロスして炎上する。
—
4. プロダクション運用のためのアーキテクチャ上の注意点
最後に、私がコードレビューで必ず指摘するRedis Stream特有の罠とベストプラクティスを共有する。
1. ストリームの肥大化を防ぐ(メモリ管理)
Redisはインメモリデータベースだ。Streamにメッセージを無限に溜め込むと、あっという間にメモリを食いつぶし、OOMを引き起こすか、maxmemoryポリシーによって古いデータが勝手に吹き飛ぶ。
メッセージを書き込む(`XADD`)際は、必ず上限を設けるか、定期的なトリミングを行うこと。
お勧め:MAXLEN と 〜(近似トリミング)を用いた安全な書き込み
約 100,000 件を維持しつつ、パフォーマンス劣化を防ぐ
127.0.0.1:6379> XADD orders MAXLEN ~ 100000 order_id 98765 amount 1500
2. コンシューマーのオートスケーリングとIDの競合
`XREADGROUP` で指定するコンシューマー名(例: `worker_1`)は、クラスター内で一意でなければならない。KubernetesのPod名やコンテナのUUIDなどを動的に割り当て、起動ごとにユニークな名前を使わせること。同じ名前のコンシューマーが複数存在すると、PELの所有権が混濁し、予期せぬ挙動を引き起こす。
3. at-least-once(最低1回は処理する)の代償と冪等性(Idempotency)
コンシューマーグループを用いたメッセージングは基本的に At-Least-Once(最低1回配送) を保証する。つまり、ネットワーク切断やワーカークラッシュによって、「実際には処理成功したが、XACKを送る前に死んだ」というケースで、同じメッセージが再配送(XCLAIMまたは再起動後の再読み込み)される可能性がある。
したがって、コンシューマー側のビジネスロジックは必ず「冪等(Idempotent)」に設計しなければならない。
例えば、決済処理であれば `order_id` をキーにしてDB側でユニーク制約を張るか、Redis側で処理済みフラグを原子的にチェックする仕組み(SETNXなど)を併用すること。
—
チーフアーキテクトからの総括
Redis Streamのコンシューマーグループは、正しく設計すればKafka並みの堅牢性とRedisならではの超低レイテンシを両立できる素晴らしい武器だ。
しかし、「とりあえず動く」コードを書くことと、「障害時にデータが消えず、自己修復するシステム」を作ることは全く次元が違う。
次に君たちが書く設計書やコードには、「XACKの漏れ対策(PEL監視とXCLAIM)」と「コンシューマーの冪等性」が当然のように組み込まれていることを期待する。
妥協のないコードを書き続けろ。
コメント