适用场景
本文适用于使用 Redis Stream 消费组承载异步任务、事件通知或轻量消息队列的系统。典型架构是生产者通过 XADD 写入 Stream,多个消费者用 XREADGROUP 拉取消息,业务成功后再执行 XACK。
当某个消费者在处理期间崩溃、被强制重启或网络中断时,已经投递但尚未确认的消息不会重新出现在 > 新消息读取结果中,而会留在消费组的 Pending Entries List(PEL)里。如果系统只读取新消息,没有故障接管机制,就会表现为“生产正常、消费者也在线,但部分任务永远不再执行”。
现象描述
常见现象包括:
- Stream 长度持续增长,但业务完成量不再同步增长;
- 消费者进程没有明显报错,读取
>时甚至一直阻塞; XINFO GROUPS中pending数量持续上升;- 某个已经离线的消费者仍持有大量 Pending 消息;
- 重启健康消费者后问题依旧存在,因为重启不会自动转移 PEL 所有权。
先确认 Stream 和消费组的整体状态:
redis-cli XINFO STREAM orders:events
redis-cli XINFO GROUPS orders:events
redis-cli XINFO CONSUMERS orders:events order-workers
重点关注:
pending:已投递但尚未XACK的数量;lag:消费组尚未投递的新消息数量;idle:消费者距离上次交互的毫秒数;inactive:消费者距离上次成功投递消息的毫秒数,较新的 Redis 版本才提供。
如果 lag 接近 0 而 pending 很高,说明新消息已经被投递,瓶颈在处理或确认阶段;如果两者都很高,则接收和处理能力可能同时不足。
可能原因
- 消费者处理完业务后、执行
XACK前崩溃,消息留在 PEL。 - 业务异常路径漏掉确认或重试处理,消息一直由原消费者持有。
- 消费者名称每次启动都随机变化,旧名称对应的 Pending 消息无人接管。
- 单条任务执行过慢,误判为故障后被其他消费者重复接管。
- 接管程序只查看前几条 Pending 消息,没有持续扫描游标。
- 先
XACK后执行业务,虽然看不到 Pending 积压,却会在进程崩溃时直接丢任务。
排查思路
1. 查看 Pending 概览
redis-cli XPENDING orders:events order-workers
返回值依次包含 Pending 总数、最小消息 ID、最大消息 ID,以及各消费者持有数量。先用它判断积压是否集中在少数消费者。
2. 展开 Pending 明细
redis-cli XPENDING orders:events order-workers - + 20
每条结果包含消息 ID、当前所有者、空闲时间和投递次数。空闲时间很长且所有者已经离线,是适合接管的候选;投递次数持续升高,则更像业务永久失败或毒消息,不能无限重试。
也可以只查看指定消费者:
redis-cli XPENDING orders:events order-workers - + 20 worker-3
3. 确认原消费者是否真的失效
不要只根据一次 idle 值立即接管。应同时检查:
- 实例是否仍在运行,是否正在执行长任务;
- 最近业务日志中是否已经开始处理该消息;
- 正常任务处理耗时的 P99;
- 接管阈值是否明显高于最大合理处理时长。
例如业务 P99 为 20 秒、极端任务不超过 60 秒,可以先将接管阈值设为 120 秒,并结合监控逐步调整。阈值过短会让正常任务被并发执行两次。
4. 读取原消息内容
redis-cli XRANGE orders:events 1726123456789-0 1726123456789-0
如果查不到消息,而 PEL 中仍有对应 ID,需检查是否过早执行了 XTRIM。删除 Stream 中的实体并不会自动清理所有历史 PEL 引用,接管时可能遇到消息内容为空的情况。
安全接管方案
Redis 6.2 及以上优先使用 XAUTOCLAIM。它会扫描 PEL,把空闲时间超过阈值的消息转移给当前消费者,并返回下一次扫描游标。
redis-cli XAUTOCLAIM \
orders:events \
order-workers \
worker-recovery-1 \
120000 \
0-0 \
COUNT 100
参数含义:
orders:events:Stream 键;order-workers:消费组;worker-recovery-1:新的消息所有者;120000:仅接管空闲至少 120 秒的消息;0-0:从 PEL 起点开始扫描;COUNT 100:限制单批扫描量,避免一次处理过多。
返回的第一个值是下一次扫描游标。必须持续用该游标调用,直到游标回到 0-0;不能假设一次调用会返回所有候选消息。
Python 恢复循环示例
下面示例使用 redis-py,只在业务处理成功后确认消息,并为永久失败设置投递次数上限。业务处理必须具备幂等性。
from __future__ import annotations
import logging
from collections.abc import Mapping
import redis
STREAM = "orders:events"
GROUP = "order-workers"
CONSUMER = "worker-recovery-1"
MIN_IDLE_MS = 120_000
MAX_DELIVERIES = 5
logger = logging.getLogger(__name__)
client = redis.Redis(host="127.0.0.1", port=6379, decode_responses=True)
def process_order(message_id: str, fields: Mapping[str, str]) -> None:
"""按业务唯一键幂等处理订单事件。"""
order_id = fields["order_id"]
# 示例:数据库写入应以 order_id + event_type 建立唯一约束。
logger.info("处理订单事件", extra={"message_id": message_id, "order_id": order_id})
def delivery_count(message_id: str) -> int:
"""读取单条 Pending 消息当前的投递次数。"""
rows = client.xpending_range(STREAM, GROUP, min=message_id, max=message_id, count=1)
if not rows:
return 0
return int(rows[0]["times_delivered"])
def recover_pending() -> None:
"""分批接管超时消息,直到完成一轮 PEL 扫描。"""
cursor = "0-0"
while True:
cursor, messages, _deleted_ids = client.xautoclaim(
STREAM,
GROUP,
CONSUMER,
min_idle_time=MIN_IDLE_MS,
start_id=cursor,
count=100,
)
for message_id, fields in messages:
attempts = delivery_count(message_id)
if attempts > MAX_DELIVERIES:
client.xadd(
f"{STREAM}:dead",
{**fields, "source_message_id": message_id, "failure_reason": "max_deliveries"},
)
client.xack(STREAM, GROUP, message_id)
logger.error(
"消息超过最大投递次数,已转入死信流",
extra={"message_id": message_id, "delivery_count": attempts},
)
continue
try:
process_order(message_id, fields)
except (KeyError, ValueError, redis.RedisError) as exc:
logger.warning(
"接管消息处理失败,保留在 Pending 中等待重试",
extra={"message_id": message_id, "delivery_count": attempts, "error": str(exc)},
)
continue
client.xack(STREAM, GROUP, message_id)
if cursor == "0-0":
break
if __name__ == "__main__":
recover_pending()
不同 redis-py 版本对 xautoclaim 返回值的封装可能不同,上线前应在当前依赖版本中做集成测试。死信转移和原消息确认不是一个原子操作;若需要严格保证,可用 Lua 脚本封装,或者让死信侧也按 source_message_id 去重。
完整消费流程
正常消费者仍应读取新消息:
redis-cli XREADGROUP \
GROUP order-workers worker-1 \
COUNT 10 BLOCK 5000 \
STREAMS orders:events '>'
推荐流程如下:
- 用
XREADGROUP ... >获取从未投递的新消息; - 依据业务唯一键执行幂等处理;
- 业务提交成功后执行
XACK; - 独立恢复任务周期性执行
XAUTOCLAIM; - 超过最大投递次数的消息进入死信 Stream,并触发告警;
- 死信修复后通过受控工具重新投递,不直接人工修改生产 PEL。
这是一种“至少一次”交付模型。Redis 可以帮助重新投递,但无法保证业务副作用只发生一次,所以数据库唯一约束、幂等表或状态机校验仍是必要条件。
容量与裁剪注意事项
不要只按固定长度激进裁剪 Stream:
redis-cli XTRIM orders:events MAXLEN ~ 100000
近似裁剪可以控制内存,但长度阈值必须覆盖业务峰值、最长故障恢复窗口和审计需求。更稳妥的做法是监控消费组的最小未确认 ID,在确认消息不再需要恢复后再制定保留策略。执行任何批量删除前,先验证 PEL 和下游持久化状态。
监控与预防
建议至少建立以下指标和告警:
- 消费组
pending总数及增长速率; - Pending 最大空闲时间;
lag及持续时间;- 单个消费者持有的 Pending 数量;
- 消息投递次数分布和死信新增速率;
- 从生产到业务完成的端到端延迟;
- 恢复任务每轮接管数、成功数、失败数与耗时。
同时落实以下工程约束:
- 消费者名称稳定且可定位到实例,不使用无法追踪的随机字符串;
- 业务处理必须幂等,确认动作必须晚于业务提交;
- 接管阈值基于真实耗时分位数,而不是凭经验随意设置;
- 恢复程序限制批量大小并有独立并发上限,避免恢复风暴;
- 为毒消息建立死信、告警和人工处置闭环;
- 升级 Redis 或客户端库时,对 PEL 扫描和返回值兼容性做回归测试。
总结
Redis Stream 中“消息不再推进”通常不是消息消失,而是消息已经进入消费组 PEL,却没有被原消费者确认。排查时应区分 lag 与 pending,再结合消费者存活状态、空闲时间和投递次数判断是否接管。使用 XAUTOCLAIM 分页扫描、设置合理空闲阈值、只在业务成功后 XACK,并用幂等和死信机制承接至少一次交付语义,才能让故障恢复既不漏消息,也不因重复执行制造新的事故。
Discussion
评论