Redis Pub/Subの深層:SUBSCRIBEが引き起こす「ブロッキングの罠」と、プロが選ぶ堅牢な設計パターン
こんにちは。アーキテクトの私だ。
今日のコードレビューで、あるジュニアエンジニアがこんなコード書いてきた。
「非同期ワーカーの通知用に、とりあえずRedisの `SUBSCRIBE` を使って常時待ち受ける実装にしました!」
……待て、手を止めろ。その設計、本番環境で確実にシステムを沈没させる。
RedisのPub/Sub機能、特に `SUBSCRIBE` コマンドは、一見すると軽量でエレガントなメッセージングに見える。だが、その内部挙動と「ブロッキング状態」の本質を理解していない者が手を出すと、クラスタ全体を巻き込む障害の引き金になる。
今日は、Redisのデータ構造とクライアントのライフサイクルを知り尽くした我々が、`SUBSCRIBE` の正体を丸裸にし、実務で絶対に破綻しないアーキテクチャを伝授しよう。
—
1. `SUBSCRIBE` の正体:なぜクライアントはブロックされるのか?
まず、Redisというデータベースの根本思想を思い出してほしい。Redisは基本的に「リクエスト・レスポンス」型のシングルスレッド(I/O多重化)サーバーだ。クライアントがコマンドを送り、サーバーが即座に結果を返す。
しかし、`SUBSCRIBE` を実行した瞬間、そのコネクションの運命は劇的に変わる。
状態遷移:通常モードから「Pub/Sub専用モード」へ
クライアントが `SUBSCRIBE channel_alpha` を発行すると、Redisサーバーはそのクライアント専用のコネクションをPub/Subモードに遷移させる。
このモードに入った瞬間、コネクションの挙動は一変する。
- 受け付けられるコマンドの激減: `PING`, `SUBSCRIBE`, `UNSUBSCRIBE`, `PSUBSCRIBE`, `PUNSUBSCRIBE` など、購読管理に関する一部のコマンドしか受け付けなくなる。
- 通常のデータ操作の拒絶: 例えば、同じコネクションで `GET` や `SET` を叩こうものなら、Redisは容赦なくエラーを返す。
- ブロッキングの発生: 購読したチャネルにメッセージが流れてくるまで、そのコネクションのI/Oスレッドはブロック(待機状態)される。
【重要】スレッドとメモリの物理的現実
ここで勘違いしてはいけないのが、「Redisサーバー側がスレッドを専有してブロックしているわけではない」という点だ。Redisのコアはシングルスレッドでイベントループを回している。
問題はクライアント側(アプリケーション側)のコネクションだ。
`SUBSCRIBE` を発行したスレッド(あるいはプロセス)は、ソケットからの読み込み(`read()` システムコール)でブロックされる。つまり、1つのコネクションを `SUBSCRIBE` に占有された場合、そのアプリケーションプロセスは、メッセージが飛んでくるまで他の仕事がそのコネクション上で一切できなくなる。
マルチスレッド言語(Java, Go, C#など)であれば専用のコネクションプールからスレッドを切り出せばいいが、シングルスレッドベースのランタイム(Node.jsやPythonの同期コードなど)でこれをやると、アプリケーション全体のライフラインが断たれる。
—
2. 実務で直面する「3つの地雷」
私が過去の障害対応で目撃した、Pub/Sub起因のトラブルのほとんどは以下の3点に集約される。
地雷1:バッファ溢れ(Client Output Buffer Limits)
もし、メッセージの生産速度(Publisher)が、購読側の処理速度(Subscriber)を大幅に上回ったらどうなるか?
Redisサーバーは、送信しきれないメッセージをクライアントごとの出力バッファ(Output Buffer)に溜め込む。
これが限界値(`client-output-buffer-limit pubsub` で設定される閾値)を超えると、Redisはどうするか。
「容赦なくそのクライアントのコネクションを切断する」。
切断された側は、再接続と再購読のロジックが甘いと、メッセージのロスト(消失)祭りだ。Pub/Subはメッセージの永続化(Persistence)を一切行わない「火縄銃」のような仕組みであることを忘れてはならない。
地雷2:コネクションの塩漬けとスケーラビリティの限界
マイクロサービスにおいて、10個のサービスがそれぞれ5つのイベントチャネルを購読し始めたとする。あっという間に数十、数百の「常時ブロックされたコネクション」がRedisサーバーに張り付くことになる。
Redisは数万のコネクションをさばけるが、Pub/Subモードのコネクションが増えると、ブロードキャスト発生時のCPU負荷(全購読者へのパケット送信処理)が急増する。安易なPub/Subの多用は、RedisをCPUバウンドなボトルネックに変えてしまう。
地雷3:リコネクト時のメッセージロスト
ネットワークの瞬断やRedisのフェイルオーバー(Master/Replicaの切り替わり)が発生した瞬間、`SUBSCRIBE` コマンドのコネクションは切断される。
切断されている間に流れたメッセージは、二度と戻らない。 RedisのPub/Subには「過去のメッセージをキャッチアップする」機能(オフセット概念)が一切ないからだ。
—
3. 実践:堅牢なPub/Sub設計パターン
では、我々プロのエンジニアはこの制約の上でどうシステムを設計すべきか。
現場で使える2つのアプローチを提示しよう。
パターンA:完全分離型・専用ワーカーアーキテクチャ
もし `SUBSCRIBE` を使うなら、「メインのWebアプリケーションプロセスとは完全に切り離された、専用の常駐ワーカープロセス」として設計しなけらばならない。
以下は、Python(`redis-py`)を用いた、堅牢なサブスクライバーワーカーの骨組みだ。
import redis
import time
import logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(“RedisWorker”)
REDIS_HOST = “localhost”
REDIS_PORT = 6379
CHANNEL_NAME = “events:system-alerts”
def create_redis_connection():
# タイムアウト等を厳格に設定した専用クライアント
return redis.Redis(host=REDIS_HOST, port=REDIS_PORT, decode_responses=True)
def resilient_subscriber():
“””
ネットワーク切断やRedisの再起動に耐える、堅牢なサブスクライバー
“””
while True:
client = create_redis_connection()
pubsub = client.pubsub()
try:
# チャネルの購読
pubsub.subscribe(CHANNEL_NAME)
logger.info(f”Successfully subscribed to {CHANNEL_NAME}. Waiting for messages…”)
# listen() はブロッキングイテレータ。ここでメッセージを待ち受ける。
for message in pubsub.listen():
# pubsubの初期化時に購読確認メッセージが流れるため、typeを検証する
if message[‘type’] == ‘message’:
handle_message(message[‘data’])
except redis.ConnectionError as e:
logger.error(f”Redis connection lost: {e}. Reconnecting in 5 seconds…”)
time.sleep(5)
except Exception as e:
logger.critical(f”Unexpected error in pubsub loop: {e}”)
time.sleep(5)
finally:
try:
pubsub.close()
except:
pass
def handle_message(data):
# ここで重い処理や外部APIコールをしてはならない。
# キューイングするか、別スレッドに処理を逃がすべき。
logger.info(f”Received event data: {data}”)
if __name__ == “__main__”:
resilient_subscriber()
このコードの設計ポイント
1. メインアプリからの隔離: このスクリプトはHTTPリクエストを処理するAPIサーバーとは別プロセスで動かす。
2. 自動リカバリ(Resilience): `redis.ConnectionError` をキャッチし、無限ループとバックオフ(待機時間)を挟んで自動で再接続する。
3. ライフサイクル管理: `finally` ブロックで確実に古いPub/Subオブジェクトを破棄し、ソケットリークを防ぐ。
—
パターンB:【最高峰の選択】Pub/Subを捨てて「Streams」を使うべきケース
もしあなたが、「メッセージが絶対にロストしては困る」「コンシューマーが落ちても、復帰後に未処理のメッセージを処理したい(At-least-once delivery)」という要件を抱えているなら、`SUBSCRIBE` を使うのは今すぐやめろ。
Redis 5.0以降で導入された `Redis Streams` (`XADD`, `XREADGROUP`) を使うべきだ。
Streamsは、Pub/Subの「リアルタイム配信」の特性を持ちながら、以下の強力なモダン機能を持つ。
- 永続化: メッセージがメモリ(およびAOF/RDB)に蓄積される。
- コンシューマーグループ: 複数ワーカーで負荷分散しつつ、誰がどのメッセージを処理したか(ACK管理)を追跡できる。
- ペンディング管理(PEL): ワーカーがクラッシュしても、未処理メッセージを別のワーカーが再取得(Claim)できる。
Pub/Sub vs Streams のアーキテクチャ比較
| 評価軸 | Pub/Sub (`SUBSCRIBE`) | Redis Streams (`XREADGROUP`) |
| :— | :— | :— |
| メッセージ永続化 | なし(その瞬間に捨てられる) | あり(ストレージに保持される) |
| 配信保証 | Fire-and-forget(At-most-once) | 確実な配信(At-least-once) |
| スケーリング | 全員が同じメッセージを受信(Pub) | グループ内で分散処理(Consumer Group) |
| 障害耐性 | 接続断時のロストあり | 復帰後のリトライ・未処理回収可能 |
| ユースケース | ライブチャット、簡易なステータス同期 | 堅牢なジョブキュー、イベント駆動アーキテクチャ |
—
4. チーフアーキテクトからの最終提言
Redisの `SUBSCRIBE` は、そのシンプルさゆえに魔力を持っている。
「とりあえずリアルタイム通知が欲しいからPub/Subで」という設計は、プロトタイプや社内ツールであれば許されるかもしれない。しかし、プロダクションレベルのシステムにおいては、「メッセージのロストがビジネスにどのような影響を与えるか」を常に逆算して選定しなければならない。
- 揮発的で、消えても困らない「純粋な通知(例: UIの再描画トリガー)」 ➡ `SUBSCRIBE` を専用プロセスで使うのはアリ。
- データとしての価値があり、ロストが許されない「イベント(例: 注文確定、決済通知)」 ➡ `SUBSCRIBE` は厳禁。`Redis Streams` か、Kafka/RabbitMQなどの本格的なメッセージbrokerを採用せよ。
道具の特性を正しく理解し、適材適所で使い倒すこと。それこそが、我々エンジニアが守るべき最高のアーキテクチャ原則だ。コードレビューの基準を、今日から一段階引き上げよう。
コメント