【実務・中級編】 Streamメッセージ承認と管理 – Redis

【Redis Streams 徹底解説】`XACK`・`XPENDING`・`XCLAIM` が織りなす「絶対にメッセージをロストしない」堅牢な非同期処理設計

テックリードの私だ。コードレビューや設計レビューで、こんな設計に遭遇したことはないか?

> 「とりあえず Redis Streams にメッセージをブチ込んで、コンシューマーが `XREADGROUP` で読んだら処理完了とするよ」

―― 待て、正気か?

その設計のままプロダクション環境にリリースしたら、数日後に「高負荷時にプロセスが死んだらデータが消えた」「メッセージが宙に浮いたまま二度と処理されない」というインシデントを踏み抜くことになる。Redis は高速だが、デフォルトでは「アット・リースト・オンズ(At-least-once: 少なくとも1回は配送)」を保証しない。

Redis Streams を真に「信頼できるメッセージング基盤」として昇華させるためには、コンシューード(読み出し)の先にある、ACK(承認)とフェイルオーバー(処理権限の移譲)のライフサイクルを完璧に制御しなければならない。

今回は、実務の現場で絶対に知っておくべき `XACK`、`XPENDING`、`XCLAIM` の核心と、それらを組み合わせた「死なないワーカーの設計パターン」を伝授しよう。

—

1. 復習:Consumer Group におけるメッセージの「暗黒面」

Consumer Group(`XREADGROUP`)を使うと、複数のワーカー間でメッセージを分散処理できる。ここでエンジニアが勘違いしやすい最大の罠がこれだ。

`XREADGROUP GROUP g1 c1 BLOCK 0 STREAMS mystream >`

このコマンドを発行した瞬間、Redis 内部で何が起きているか?
Redis は、そのメッセージをコンシューマー `c1` の「PEL(Pending Entries List:未了エントリリスト)」に登録する。

つまり、メッセージは「コンシューマーに渡された」が、「正常に処理が完了した(ACKされた)」わけではない。コンシューマー `c1` が処理の途中で OOM Killer に殺されたり、ネットワークが切断されたりした場合、そのメッセージは `c1` の PEL の中に永遠に幽霊のように残り続ける。

この幽霊たちを成仏させ、再処理するための三種の神器が `XACK`, `XPENDING`, `XCLAIM` である。

—

2. 三種の神器のメカニズムと実務的解釈

① `XACK`: 処理完了の狼煙(のろし)

  • 役割: PEL からメッセージを削除し、「このメッセージのライフサイクルは完了した」と Redis に宣言する。
  • 鉄則: ビジネスロジックが完全に成功した「最後」にのみ呼ぶこと。 例外キャッチの `catch` や `finally` で安易に呼ぶな。

構文: XACK <ストリーム名> <グループ名> <メッセージID>
XACK mystream g1 1689000000000-0

② `XPENDING`: ゾンビメッセージの観測所

  • 役割: PEL に溜まっている「まだ ACK されていないメッセージ」をスキャンする。
  • 実務での使いどころ: どのコンシューマーがサボっているか(あるいは死んでいるか)、どれくらい処理遅延(Idle Time)が発生しているかを監視するために使う。

グループ全体の未処理状況のサマリーを取得
XPENDING mystream g1
実行結果例:
1) (integer) 3 # 未処理の総数
2) “1689000000000-0” # 最古のメッセージID
3) “1689000001000-0” # 最新のメッセージID
4) 1) 2) “worker-A” # コンシューマーごとの未処理数 (“worker-A” が 3件持っている)
(integer) 3

詳細な明細を見るには `XPENDING` に範囲を指定する。

アイドル時間が長い顺に詳細を見る(後述の XCLAIM の前段階で使う)
XPENDING mystream g1 – + 10

③ `XCLAIM`: 権利剥奪と強奪(フェイルオーバー)

  • 役割: 死んだ(あるいは処理が異常に遅い)コンシューマーの PEL から、メッセージの所有権を別のコンシューマーに強制移譲する。
  • 肝: `MIN-IDLE-TIME`(ミリ秒)を指定することで、「指定時間以上放置されているメッセージだけを奪い取る」という安全装置が働く。

worker-A が死んでいると仮定し、worker-B が 60秒以上放置されたメッセージを奪う
XCLAIM <ストリーム> <グループ> <新オーナー> <最小アイドル時間(ms)> …
XCLAIM mystream g1 worker-B 60000 1689000000000-0

—

3. 【設計パターン】死なないコンシューマー・ガーディアンの実装

では、これらを組み合わせて「如何なる障害(プロセス強制終了、ネットワーク分断)があってもメッセージを取りこぼさない」堅牢なワーカーのアーキテクチャをコードで示そう。

実務では、メインのメッセージ処理ループとは別に、「リクレイマー(回収屋)スレッド / プロセス」を常駐させるのが定石だ。

アーキテクチャ概要

1. ワーカー・ループ: `XREADGROUP` で取得 ➔ 処理 ➔ 成功時のみ `XACK`
2. リクレイマー・ループ: 定期的に `XPENDING` で古い PEL をスキャン ➔ 放置されていたら `XCLAIM` で自衛(または別ワーカー)に奪取して再処理。

Python による堅牢なワーカーの実装例

import time
import redis

client = redis.Redis(host=”localhost”, port=6379, decode_responses=True)

STREAM_NAME = “orders”
GROUP_NAME = “order-processors”
CONSUMER_NAME = “worker-1″
IDLE_TIME_THRESHOLD_MS = 60000 # 60秒放置されたらスタックとみなす

def setup_stream():
try:
client.xgroup_create(STREAM_NAME, GROUP_NAME, id=”0”, mkstream=True)
except redis.exceptions.ResponseError:
pass # すでにグループが存在する場合は無視

def process_message(message_id, data):
“””ビジネスロジック。ここで例外が発生すれば XACK は呼ばれない”””
print(f”Processing message {message_id}: {data}”)
# 例: 外部API呼び出しやDB書き込み
if “fail” in data.get(“action”, “”):
raise RuntimeError(“Intentional processing failure”)

def worker_loop():
setup_stream()
while True:
try:
# 1. 新規メッセージの取得 (自分宛ての未読 > を取得)
streams = client.xreadgroup(
GROUP_NAME, CONSUMER_NAME, {STREAM_NAME: “>”}, count=1, block=2000
)

if not streams:
# メッセージがない暇な時間に、放置されたメッセージの回収(Claim)を行う
reclaim_abandoned_messages()
continue

for stream, messages in streams:
for message_id, data in messages:
try:
# 2. 処理実行
process_message(message_id, data)

# 3. 正常終了時のみ ACK
client.xack(STREAM_NAME, GROUP_NAME, message_id)
print(f”Successfully acknowledged: {message_id}”)

except Exception as e:
print(f”Error processing {message_id}: {e}”)
# ACKしない。これによりメッセージはPELに残る。

except Exception as e:
print(f”Worker loop error: {e}”)
time.sleep(1)

def reclaim_abandoned_messages():
“””死んだコンシューマーの PEL をスキャンし、自分のものに書き換えて再処理する”””
try:
# PELから古いエントリを検出し、アイドル時間が閾値を超えているものを取得
# XPENDINGのレンジ取得: min, max, count
pendings = client.xpending_range(
STREAM_NAME, GROUP_NAME, min=”-“, max=”+”, count=10
)

for p in pendings:
message_id = p[“message_id”]
idle_time = p[“time_since_delivered”]
current_owner = p[“consumer”]

if current_owner == CONSUMER_NAME:
continue # 自分が持っているものはスキップ(自分の通常ループで処理するか別途ハンドリング)

if idle_time > IDLE_TIME_THRESHOLD_MS:
print(f”Claiming abandoned message {message_id} from {current_owner} (Idle: {idle_time}ms)”)

# XCLAIM で強制的に自分のものにする
claimed = client.xclaim(
STREAM_NAME, GROUP_NAME, CONSUMER_NAME,
min_idle_time=IDLE_TIME_THRESHOLD_MS,
message_ids=[message_id]
)

if claimed:
for message_id, data in claimed:
# 奪取したメッセージを処理して ACK
try:
process_message(message_id, data)
client.xack(STREAM_NAME, GROUP_NAME, message_id)
print(f”Successfully recovered & acknowledged: {message_id}”)
except Exception as e:
print(f”Failed to process recovered message {message_id}: {e}”)
# ここでデッドレターキュー(DLQ)への退避やリトライ回数制限のインクリメントが必要
handle_dead_letter(message_id, data)

except Exception as e:
print(f”Reclaim error: {e}”)

def handle_dead_letter(message_id, data):
“””何回も失敗するメッセージの無限ループ(毒薬メッセージ)を防ぐための退避処理”””
# 実際には別の Stream や DB に移して XACK する
print(f”Moving {message_id} to Dead Letter Queue.”)
client.xack(STREAM_NAME, GROUP_NAME, message_id)

if __name__ == “__main__”:
worker_loop()

—

ジックの裏側には、こうした堅牢性の担保が不可欠だ。

—

4. チーフアーキテクトが警鐘を鳴らす「パフォーマンスと運用の罠」

最後に、プロダクション環境でこの設計を運用する際に、必ず直面する「落とし穴」をシェアしよう。

1. 無限リトライ(Poison Pill)の恐怖

  • `XCLAIM` で無限にメッセージを回収し続けると、バグを含んだメッセージ(ポイズンピル)が原因で、ワーカーが永遠にエラーとリカバリを繰り返す。
  • 対策: メッセージのメタデータ(あるいはRedis上の別ハッシュ)に「配信試行回数(Delivery Counter)」を持たせ、3回失敗したら強制的に `XACK` してデッドレター・ストリーム(DLQ)へ退避させるロジックを必ず入れろ。

2. PEL の肥大化によるメモリ圧迫

  • ACK をサボり続けると、Redis の RAM 上にある PEL が肥大化し、メモリを圧迫する。
  • PEL は Redis のメモリ内に保持されるため、コンシューマーがゴミデータを残したまま消滅し続けるとメモリリーク的な状況に陥る。モニタリングで `XPENDING` の件数を常に監視せよ。

3. クラスタ環境における `XCLAIM` の注意点

  • Redis Cluster を使っている場合、Stream はハッシュスロット単位で分割される。`XCLAIM` や `XACK` は同一のスロット(同じノード)をターゲットにする必要があるため、キーのルーティング設計(ハッシュタグ `{mystream}` の利用など)に注意すること。

まとめ

Redis Streams は、単なる「速いキュー」ではない。
`XREADGROUP` で取り出し、ビジネスロジックを完遂し、`XACK` で手放す。そして、脱落したワーカーの遺産は `XPENDING` で見つけ出し、`XCLAIM` で回収する。

このライフサイクルをコードに組み込んで初めて、あなたは「プロフェッショナルなバックエンドエンジニア」と名乗る資格を得る。
次回の設計レビューでは、コンシューマーの死活監視と ACK のタイミングについて、チームメンバーに厳しい視線を向けてやってほしい。

コメント

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