Redis Stream作为5.0版本引入的数据结构,凭借消费者组机制成为轻量级消息队列的热门选择。但不少团队在上线后都遇到过同一个问题:生产消息明明已经写入,消费者日志也显示处理完成,队列的Pending Entries List(PEL)却越积越多,内存占用持续上涨。要定位这类问题,第一步就是用好XPENDING命令。它就像一个透视镜,能把你关心的待处理消息一条条翻出来看清楚。

Pending列表是怎么产生的
要理解XPENDING的输出,先得明白Pending Entries List的机制。当消费者通过XREADGROUP读取消息时,Redis并不会直接认为消息已被处理,而是在该消费者组关联的PEL结构中登记一条记录,包含消息ID、消费者名称、投递次数和上次投递时间戳。只有消费者显式执行XACK后,这条记录才会从PEL中移除。换句话说,PEL本质上是一张“已投递但未确认”的账本。
这张账本会膨胀的场景主要有三类:第一,消费者代码忘记调用XACK,属于最常见的编码疏漏;第二,消费者进程崩溃或被杀掉,读取了消息却没来得及确认,消息就永远挂在那个消费者名下;第三,处理逻辑抛出异常,消息确认被跳过,反复重投后仍然失败。无论哪种情况,你都需要先通过XPENDING弄清楚积压的规模和分布,再决定用XCLAIM转移还是做其他补偿。
可以用一个简单实验观察PEL的形成过程:
# 创建Stream和消费者组 XADD mystream * task order-1001 XGROUP CREATE mystream mygroup 0 # 消费者读取但不确认 XREADGROUP GROUP mygroup consumer-1 COUNT 1 STREAMS mystream > # 查看Pending列表,此时应有1条记录 XPENDING mystream mygroup
XPENDING基础查询与返回字段详解
XPENDING最简形式只需要两个参数:Stream键名和消费者组名。执行后返回的是汇总信息,包含五个部分:PEL中的消息总数、消息ID的最小值和最大值、以及每个消费者名下的待处理消息数量分布。这个汇总视图非常适合做监控指标,比如接入Prometheus后可以直观看到积压趋势。
# 汇总查询 XPENDING mystream mygroup # 返回示例 # 1) (integer) 3 -- 待处理消息总数 # 2) "1718000000000-0" -- 最小消息ID # 3) "1718000005000-0" -- 最大消息ID # 4) 1) 1) "consumer-1" # 2) (integer) 2 -- consumer-1持有2条 # 2) 1) "consumer-2" # 2) (integer) 1 -- consumer-2持有1条
如果需要看每条消息的明细,就得使用扩展格式,在命令后面追加START、END、COUNT三个参数。这时每条记录会展示四个关键信息:消息ID、当前持有该消息的消费者名、该消息的投递次数delivery counter、以及以毫秒计的空闲时长idle time。投递次数尤其值得重点关注,如果某条消息的counter值很高,通常意味着它反复被投递又反复失败,大概率是处理逻辑存在bug,或者消息本身是脏数据。
# 查询前10条明细,ID范围为全部 XPENDING mystream mygroup - + 10 # 返回示例 # 1) 1) "1718000000000-0" # 2) "consumer-1" # 3) (integer) 5 -- 已投递5次 # 4) (integer) 120000 -- 空闲120秒
此外扩展格式还支持两个可选参数。IDLE可以过滤出空闲时间超过指定毫秒数的消息,非常适合找出“卡死”的任务;consumer参数则限定只查某个消费者的名下消息,排查单个实例的故障时非常方便。例如XPENDING mystream mygroup - + 10 IDLE 60000 consumer-1表示只查consumer-1名下空闲超过1分钟的最多10条消息。
游标分页遍历大规模Pending消息
当积压量达到几万甚至几十万条时,一次性- + COUNT的大范围查询会给Redis主线程带来压力,也会让客户端一次收到过大的响应包。Redis 6.2开始为XPENDING引入了游标机制,写法是在COUNT后面加上大写的关键字和游标值。首次查询游标传0,Redis返回结果的同时会给出下一页的游标位置,拿到空游标则表示遍历结束。
import redis
r = redis.Redis(host='127.0.0.1', port=6379, decode_responses=True)
def scan_pending(stream, group, batch=100):
cursor = '0'
stuck_messages = []
while cursor != None:
# 返回值第一项为下一页游标,第二项为消息明细列表
cursor, entries = r.xpending_range(
stream, group, min='-', max='+',
count=batch, idletime=300000, consumername=None
) if False else (None, [])
# 简化演示:直接调用xpending_range
result = r.xpending_range(stream, group, min='-', max='+', count=batch)
if not result:
break
stuck_messages.extend(result)
# 以最后一条消息ID作为下一页起点
last_id = result[-1]['message_id']
# 实际游标分页需借助XPENDING ... IDLE ... ... 游标语法
break
return stuck_messages
上面演示了客户端遍历的思路,实际在Redis 6.2及以上环境更推荐直接使用带游标的命令形式:XPENDING key group IDLE idle [consumer] count ms-count,其中ms-count位置传入上一轮返回的游标。相比传统的按ID区间翻页,游标方式不需要客户端自己记录边界,而且支持与IDLE过滤组合使用,代码更简洁也更不容易漏数据。遍历拿到问题消息清单后,通常下一步就是对空闲过久的消息执行XCLAIM或XAUTOCLAIM,把它们转移到健康的消费者重新处理,处理成功后再XACK,整个故障恢复链路就闭环了。
实战:清理僵尸消费者名下的消息
假设线上某台机器宕机,它上面的消费者consumer-2名下挂着几百条未确认消息。完整处理流程是:先用XPENDING确认consumer-2名下的消息规模和空闲时间,再用XCLAIM批量把消息转移给活跃的consumer-1,并设置最小空闲时间门槛,避免误抢正在处理中的消息,最后由consumer-1处理完逐条XACK。
# 第一步:查看consumer-2名下积压情况 XPENDING mystream mygroup - + 100 consumer-2 # 第二步:把空闲超过60秒的消息转移给consumer-1 # 格式:XCLAIM key group consumer min-idle-time id [id ...] XCLAIM mystream mygroup consumer-1 60000 1718000000000-0 1718000000001-0 # 第三步:处理成功后确认 XACK mystream mygroup 1718000000000-0 # 也可以用XAUTOCLAIM自动批量领取 XAUTOCLAIM mystream mygroup consumer-1 60000 0 COUNT 50
需要注意,XCLAIM会重置消息的投递计数相关行为,可以配合JUSTID参数只转移所有权不返回消息体,减少网络开销。另一个实践建议是在业务代码里对投递次数设置上限,比如通过XPENDING明细发现某条消息delivery counter超过10次,就直接进入死信处理流程记录告警,而不是无限重试拖垮系统。把XPENDING的定期巡检做成定时任务,配合监控告警,就能在消息积压演变成事故之前及时介入,这也是使用Stream做可靠消息队列必不可少的一环。
Redis XPENDINGStream待处理消息消息队列修改时间:2026-09-05 17:58:51