适用场景

本文适用于使用 Redis Stream 消费组承载异步任务、事件通知或轻量消息队列的系统。典型架构是生产者通过 XADD 写入 Stream,多个消费者用 XREADGROUP 拉取消息,业务成功后再执行 XACK

当某个消费者在处理期间崩溃、被强制重启或网络中断时,已经投递但尚未确认的消息不会重新出现在 > 新消息读取结果中,而会留在消费组的 Pending Entries List(PEL)里。如果系统只读取新消息,没有故障接管机制,就会表现为“生产正常、消费者也在线,但部分任务永远不再执行”。

现象描述

常见现象包括:

  • Stream 长度持续增长,但业务完成量不再同步增长;
  • 消费者进程没有明显报错,读取 > 时甚至一直阻塞;
  • XINFO GROUPSpending 数量持续上升;
  • 某个已经离线的消费者仍持有大量 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 很高,说明新消息已经被投递,瓶颈在处理或确认阶段;如果两者都很高,则接收和处理能力可能同时不足。

可能原因

  1. 消费者处理完业务后、执行 XACK 前崩溃,消息留在 PEL。
  2. 业务异常路径漏掉确认或重试处理,消息一直由原消费者持有。
  3. 消费者名称每次启动都随机变化,旧名称对应的 Pending 消息无人接管。
  4. 单条任务执行过慢,误判为故障后被其他消费者重复接管。
  5. 接管程序只查看前几条 Pending 消息,没有持续扫描游标。
  6. 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 '>'

推荐流程如下:

  1. XREADGROUP ... > 获取从未投递的新消息;
  2. 依据业务唯一键执行幂等处理;
  3. 业务提交成功后执行 XACK
  4. 独立恢复任务周期性执行 XAUTOCLAIM
  5. 超过最大投递次数的消息进入死信 Stream,并触发告警;
  6. 死信修复后通过受控工具重新投递,不直接人工修改生产 PEL。

这是一种“至少一次”交付模型。Redis 可以帮助重新投递,但无法保证业务副作用只发生一次,所以数据库唯一约束、幂等表或状态机校验仍是必要条件。

容量与裁剪注意事项

不要只按固定长度激进裁剪 Stream:

redis-cli XTRIM orders:events MAXLEN ~ 100000

近似裁剪可以控制内存,但长度阈值必须覆盖业务峰值、最长故障恢复窗口和审计需求。更稳妥的做法是监控消费组的最小未确认 ID,在确认消息不再需要恢复后再制定保留策略。执行任何批量删除前,先验证 PEL 和下游持久化状态。

监控与预防

建议至少建立以下指标和告警:

  • 消费组 pending 总数及增长速率;
  • Pending 最大空闲时间;
  • lag 及持续时间;
  • 单个消费者持有的 Pending 数量;
  • 消息投递次数分布和死信新增速率;
  • 从生产到业务完成的端到端延迟;
  • 恢复任务每轮接管数、成功数、失败数与耗时。

同时落实以下工程约束:

  • 消费者名称稳定且可定位到实例,不使用无法追踪的随机字符串;
  • 业务处理必须幂等,确认动作必须晚于业务提交;
  • 接管阈值基于真实耗时分位数,而不是凭经验随意设置;
  • 恢复程序限制批量大小并有独立并发上限,避免恢复风暴;
  • 为毒消息建立死信、告警和人工处置闭环;
  • 升级 Redis 或客户端库时,对 PEL 扫描和返回值兼容性做回归测试。

总结

Redis Stream 中“消息不再推进”通常不是消息消失,而是消息已经进入消费组 PEL,却没有被原消费者确认。排查时应区分 lagpending,再结合消费者存活状态、空闲时间和投递次数判断是否接管。使用 XAUTOCLAIM 分页扫描、设置合理空闲阈值、只在业务成功后 XACK,并用幂等和死信机制承接至少一次交付语义,才能让故障恢复既不漏消息,也不因重复执行制造新的事故。