Redis Stream徹底攻略:分散システムにおける「確実なメッセージ処理」の設計論
Redisを単なる「高速なキーバリューストア」としか見ていないなら、それは大きな損失だ。
特に、Redis 5.0で導入されたStream、そしてそれを制御するコンシューマーグループは、現代の分散システムにおいて「信頼性の高いメッセージング」を構築するための最強の武器の一つだ。
今日は、小手先のコマンド解説ではなく、実務の戦場で生き残るための「Redis Streamの極限的な設計パターン」について深掘りしていく。
—
1. なぜ「Pub/Sub」ではなく「Stream」なのか?
まず、大前提を叩き込んでおく。Redisの従来のPub/Subは「火を噴いて消える花火」だ。受信側がオフラインなら、データは永遠に失われる。
一方、Streamは「永続化されたログ」だ。
コンシューマーグループ(Consumer Group)を組み合わせることで、以下の3点を保証できる。
1. 分散読み取り: 複数のコンシューマーで負荷を並列分散(メッセージの奪い合い)。
2. 到達保証: ACK(Acknowledgment)が来るまで、メッセージは「保留中(PEL)」として追跡される。
3. 耐障害性: コンシューマーが死んでも、別のコンシューマーが処理を肩代わりできる。
—
2. コンシューマーグループの「生存戦略」:設計の要諦
設計レビューで必ず聞くのは「そのメッセージ処理は冪等(べきとう)か?」だ。
Redis StreamのACKモデルは非常に強力だが、ネットワークの瞬断等で「処理は終わったがACKが届かなかった」という状況は必ず発生する。
基本の実装パターン
1. コンシューマーグループの作成
$は「現在の末尾」を指す。0なら過去の全メッセージを対象にする
XGROUP CREATE mystream mygroup $ MKSTREAM
2. メッセージの読み取り(コンシューマーAが読み取る)
> は「未配信の新しいメッセージ」を取得することを意味する
XREADGROUP GROUP mygroup consumer-a COUNT 10 STREAMS mystream >
3. 処理完了後のACK送信(これが無いとメッセージは消えない)
XACK mystream mygroup
【チーフアーキテクトの視点】
`XREADGROUP` で取得する際、IDに `>` を指定するのは基本だが、「未処理の再送(PELの回収)」には `0` を指定する。これを定期的に巡回させるプロセスを別スレッドで回すのが、堅牢なシステム設計の基本だ。
—
3. 実務でハマる「落とし穴」とパフォーマンス戦略
多くのエンジニアがここで躓く。パフォーマンスを犠牲にせず、かつデータを取りこぼさないための注意点だ。
① 「PEL(Pending Entries List)」の肥大化を許すな
`XACK` を忘れると、Redisメモリ上に「未処理リスト」が溜まり続ける。メモリ溢れを引き起こし、Redisプロセスそのものが `OOM` (Out of Memory) で落ちる。
- 対策: `XPENDING` で定期的に監視し、一定時間(例:5分以上)未処理のメッセージがあれば、処理担当を強制的に変更するか、アラートを飛ばすこと。
② メッセージIDの管理
StreamのID(例:`1620000000000-0`)はタイムスタンプベースだ。
もしアプリケーション側で独自のIDを振りたいなら、Redisの自動生成に任せるべきだ。無理にカスタムIDを入れると、Streamの順序保証が崩れ、将来的なトラブルの温床になる。
③ 「ストリームのトリミング」は必須
Streamは追加し続けるとメモリを食いつぶす。`XADD` を実行する際、`MAXLEN` オプションを必ず入れろ。
1000件を超えたら古いものを捨てる設定
XADD mystream MAXLEN ~ 1000 key value
`~` を付けることで、厳密な1000件制限ではなく「概ね1000件」とすることで、内部的なメモリ再配置のコストを抑えるのが通のやり方だ。
—
4. 堅牢な設計パターン:Workerの再起動とリカバリー
システムがクラッシュした際、誰もACKを返せない状態になる。この時、どう設計すべきか。
1. オーナーシップの移譲: `XPENDING` で古いメッセージを検知し、`XCLAIM` を使って別のコンシューマーに所有権を強制的に移譲する。
2. 冪等性の担保: 上記の再送プロセスが発生すると、同じメッセージが2回処理される可能性がある。必ずDB側で `ユニーク制約` や `処理済みフラグ` を用いた排他制御を行うこと。
—
結論:Redisは「信頼の基盤」になり得る
Redis Streamを正しく使えば、複雑なメッセージキューイングシステム(Kafkaなど)を導入する前に、多くの問題を解決できる。
- 小規模〜中規模の分散処理なら、Redis Streamで十分だ。
- ただし、ACKとPELの管理を怠れば、それは爆弾になる。
君たちのシステムが「なぜその設計なのか」を、このコードと理論で説明できるようになってほしい。Redisは、理解した者の意図を忠実に実行する、極めて優秀なエンジンだ。
さあ、コードを書こう。そして、堅牢なシステムを構築してくれ。
—
追記:もし「なぜメモリ消費が激しいのか」という質問があれば、次回は `Redisのメモリ断片化` と `Streamの内部データ構造(Radix Tree)` について語ることにしよう。
コメント