Redis Streamsの深層:なぜ「ただのログ」として使うべきではないのか
こんにちは。テックリードの私だ。
今日のコードレビューで、あるジュニアエンジニアがこんなコード書いてきた。
「監査ログの保存用に、RedisのList(`LPUSH`/`RPOP`)を使っています!」
……おいおい、ちょっと待て。現代の分散システムにおいて、イベント駆動アーキテクチャのバックボーンや、堅牢なメッセージングキューを構築する際、いつまでも古臭いListやPub/Subにしがみついている場合じゃない。Pub/Subは「撃ちっ放し(Fire-and-forget)」だからネットワーク切断でデータが消飛ぶし、Listはコンシューマーのオフセット管理を自前で実装しようとすると地獄を見る。
そこで登場するのが Redis Streams だ。
Redis 5.0で導入されたこのデータ構造は、単なる「追記型のログ」ではない。Apache KafkaのエッセンスをRedisの超高速なインメモリ空間に抽出し、洗練されたコンシューマーグループ機構まで持たせた、極めて強力なプリミティブだ。
今回は、このRedis Streamsの根幹をなす基本概念、IDの構造、そして現場で即座に使える堅牢な設計パターンを、アーキテクトの視点から叩き込む。
—
1. Streamsの正体:なぜKafkaに匹敵しうるのか
Redis Streamsは、一言で言えば「時系列でイミュータブル(不変)なエントリを追記していく、マルチコンシューマー対応のログデータ構造」だ。
従来のRedisデータ構造との違いを整理しておこう。
- Pub/Sub: データは保持されない。その瞬間繋がっているクライアントに配信されるだけ。オフセット?何それ美味しいの?の世界。
- List: 順序は保証されるが、「誰がどこまで読んだか(コンシューマーの読み込み位置)」をRedis側で覚えてくれない。アプリケーション層で死活監視やオフセット管理を実装する必要があり、バグの温床になる。
- Streams: データは永続化され(AOF/RDB)、各エントリには一意なIDが付与される。さらに、コンシューマーグループ(Consumer Groups)により、複数ワーカーでのロードバランスと「誰がどのメッセージを処理中で、どれが未完了か(ACK待ちか)」のトラッキングを完全にサポートする。
—
2. `XADD` と「時間ベースのID構造」の深層
Streamsにデータを書き込む基本コマンドが `XADD` だ。
まずはそのシグネチャと、内部で何が起きているかを見てみよう。
基本的な書き込み
127.0.0.1:6379> XADD mystream sensor_id “temp-01” value “23.5”
“1711929600000-0”
ここで注目すべきは、第2引数の “ だ。これは「RedisにエントリIDを自動生成させる」という指示である。
返却された `”1711929600000-0″` という文字列こそが、Streamsの命とも言えるIDの構造だ。
IDのフォーマット:`–`
Redis StreamsのIDは、ハイフンで区切られた2つのパーツで構成されている。
1. ミリ秒タイムスタンプ (`1711929600000`): エントリがRedisサーバーに追加された時点のUnixミリ秒(`13桁`)。これにより、Streams内のエントリは時間順に自然ソートされる。
2. シーケンス番号 (`0`): 同一ミリ秒内に複数のエントリが追加された場合、競合を防ぐためにインクリメントされるカウンター(`0`オリジン)。
💡 アーキテクトの知見:IDの自作と罠
実は、“ の代わりに手動でIDを指定することも可能だ(例: `XADD mystream 1711929600000-1 …`)。
しかし、実務において明示的なID指定は「アンチパターン」であることが多い。Redisは厳密に「直前のIDより大きいこと」を強制するため、時計の巻き戻しや順序の逆転が発生すると、容赦なく `ERR The ID specified in XADD is equal or smaller than the target stream top` というエラーを吐いて爆死する。
特別な理由(外部システムからのマイグレーション等)がない限り、ID生成は必ず “ に任せること。
—
3. 実務で直面する「時計の歪み」問題と対策
時間ベースのIDを使っている以上、避けて通れないのがOSの時刻同期(NTP)の歪みだ。
もしNTPの調整によってサーバーの時計が過去に戻った場合、何が起きるか?
Redisは「前回のIDより大きいか」をチェックしているため、時計が戻ると、新しいエントリのタイムスタンプが過去のものになり、書き込みが拒絶される。
【対策】
1. NTPのスルーモード(Slew mode)を使用する: 時刻をジャンプさせるのではなく、徐々に進める(あるいは遅らせる)設定にする。これにより、ミリ秒単位の急激な逆転を防げる。
2. 専用のコンテナ/VM環境の監視: クラウド環境(AWS EC2など)では、`chrony` などのデーモンを適切に設定し、時刻ジャンプを極力抑制する。
—
4. 堅牢な設計パターン:イベント駆動パイプラインの構築
では、実際のシステム開発でどう設計すべきか。
「Webhooksの受信バッファ」を例に、堅牢な設計パターンを示す。
アーキテクチャ概要
1. APIサーバーがWebhooksを受信し、即座に `XADD` でStreamへブロードキャスト。
2. 複数のバックグラウンドワーカー(コンシューマーグループ)が `XREADGROUP` で並行処理。
3. 処理成功したら `XACK`、失敗した(クラッシュ等)場合は未処理として保持し、後でリトライ。
1. ストリームへの書き込み(APIサーバー側)
メキシコからの決済イベントを受信したと想定
XADD payments:stream MAXLEN ~ 10000 event_id “evt_998127” amount “15000” currency “JPY”
> ⚠️ 設計上の極意 (`MAXLEN ~ 10000`):
> メモリは有限だ。無限にログが溜まり続ければ、いつかOOM(Out of Memory)でRedisがクラッシュする。`MAXLEN ~` を使うことで、おおよそ1万件にトリムしつつ、メモリ効率の良いデータ構造を維持できる(`~` は厳密な1万件ではなく、メモリ管理の効率を優先して多少前後することを許容するフラグ。実務では必須)。
2. コンシューマーグループの作成
システム起動時に、一度だけ実行する(既に存在する場合はエラーになるので注意)。
“payment-processors” というグループを作成。0 はストリームの最初から読む指定
XGROUP CREATE payments:stream payment-processors $ MKSTREAM
- `MKSTREAM`: 指定したStreamが存在しない場合に、自動的に空のStreamを作成する(Redis 6.2以降で非常に便利)。
3. ワーカーによるメッセージの読み込みと処理
ワーカーは以下のコマンドで自分の担当分のメッセージをフェッチする。
グループ “payment-processors” のコンシューマー “worker-1” が、まだ読んでいないメッセージを10件取得
XREADGROUP GROUP payment-processors worker-1 COUNT 10 BLOCK 2000 STREAMS payments:stream >
- `>`: 「他のコンシューマーにまだ配送されていない、新しいメッセージ」を指す特殊なID。
- `BLOCK 2000`: メッセージがない場合、最大2秒間ブロッキングしてCPUの無駄な高騰(ビジネスカウント)を防ぐ。
4. 処理完了の通知 (ACK)
ワーカーがメッセージの処理(DB保存や外部API呼び出しなど)を正常に終えたら、必ずACKを返す。
XACK payments:stream payment-processors 1711929600000-0
これを怠ると、メッセージは「保留リスト(Pending Entries List: PEL)」に残り続け、メモリを圧迫し続けることになる。
—
5. パフォーマンスと運用の注意点(チーフアーキテクトからの戒め)
最後に、現場で事故を起こさないための鉄則をいくつか残しておく。
1. ブロキングコマンドのタイムアウト設計:
`XREADGROUP` や `XREAD` で無限ブロック(`BLOCK 0`)を使うな。コネクションがゾンビ化し、Redisのファイル記述子(File Descriptors)を枯渇させる原因になる。必ず適切なタイムアウトを設定せよ。
2. メモリ上限とMaxmemory Policy:
Streamを使用する場合、`maxmemory` に達したときのポリシー(`noeviction` 推奨だが状況による)を誤ると、データが勝手に削除されてロストする。データロストが許されないシステムでは、`maxmemory-policy noeviction` に設定し、メモリが溢れたら書き込みを拒否するように設計すべきだ。
3. XPENDING と XCLAIM の活用:
ワーカーが処理中にクラッシュした際、PELに放置されたメッセージを救出する仕組み(Dead Letter Queue的なアプローチや、一定時間経過したメッセージを別のワーカーに再割り当てする `XCLAIM`)を必ずバッチや監視プロセスで実装しておけ。
結び
Redis Streamsは、単なる「おもちゃのキュー」ではない。適切に設計・運用すれば、大規模なトラフィックをさばく堅牢なイベントソーシング基盤として十分に機能する。
「とりあえずListで」という設計を見かけたら、この記事を突きつけてやってほしい。
「うちはRedis Streamsで、コンシューマーグループとPELまで完璧に回すんだ」とね。
あなたのアーキテクチャが、美しく、そして強靭であることを祈る。
コメント