Redis Stream是5.0版本引入的日志型数据结构,配合消费者组(Consumer Group)可以把它当成一个轻量级的消息队列来用。消费者组本质上是对Stream的一条逻辑消费链路,同一个Stream可以挂多个组,组与组之间的消费进度互相独立,组内的多个消费者则分摊消息。而这一切的起点,就是XGROUP命令族。这篇文章会详细拆解XGROUP的各个子命令,从创建到销毁,把每个参数的坑都说清楚。

一、XGROUP CREATE:创建消费者组的正确姿势
创建消费者组最基础的命令是XGROUP CREATE,它的完整语法是XGROUP CREATE key group id|$ [MKSTREAM]。其中key是Stream的键名,group是消费者组的名字,最后一个参数决定从哪个位置开始消费。0表示从头开始消费Stream中的所有历史消息,而$表示只消费创建组之后新写入的消息,这个区别非常重要,选错了要么重复处理一堆旧数据,要么漏掉本该消费的消息。
实际操作示例:
127.0.0.1:6379> XGROUP CREATE orders order-group 0 MKSTREAM OK 127.0.0.1:6379> XADD orders * event created "1718000000000-0"
上面用到了MKSTREAM选项,它的作用是:当key不存在时自动创建一个空Stream。如果不加这个选项,对一个不存在的key执行XGROUP CREATE会直接报错。这一点在设计系统启动流程时要留意——如果服务启动时Stream可能还不存在,建议始终带上MKSTREAM,避免初始化失败。
另外一个容易踩的坑是,如果指定的key存在但不是Stream类型,比如是一个List或String,命令会返回WRONGTYPE错误。所以线上操作前最好先用TYPE命令确认一下key的类型。消费者组创建成功后,可以通过XINFO GROUPS orders查看组的元信息,包括last-delivered-id、挂起的条目数等,这些指标是后续排查消费进度的基础。
二、组内消费者的增删:CREATECONSUMER与DELCONSUMER
在较老的版本中,消费者是隐式创建的,第一次执行XREADGROUP时会自动注册。但从Redis 6.2开始,提供了显式管理消费者的能力:XGROUP CREATECONSUMER key group consumer用于创建一个消费者,XGROUP DELCONSUMER key group consumer用于删除。
删除消费者时要特别注意:该消费者名下未被确认的消息不会凭空消失,而是仍然留在组的Pending Entries List(PEL)中,处于无人认领的状态。这些消息需要通过XCLAIM或XAUTOCLAIM转移给其他消费者,否则它们会一直挂在组上,造成消息堆积。一个典型的故障处理流程是这样的:
127.0.0.1:6379> XPENDING orders order-group - + 10 1) 1) "1718000000000-0" 2) "consumer-a" 3) (integer) 65 4) (integer) 1718000500000 127.0.0.1:6379> XGROUP DELCONSUMER orders order-group consumer-a (integer) 1 127.0.0.1:6379> XAUTOCLAIM orders order-group consumer-b 60000 0-0 COUNT 10
上面的命令先查看挂起的消息,发现consumer-a有一条消息超过60秒没确认,于是把它删掉,再用XAUTOCLAIM把消息转移给consumer-b重新处理。这套组合拳是处理消费者宕机后消息转移的标准做法。
值得强调的是,消费者组和消费者是两个层级的概念。组是消费进度的载体,消费者只是组内的一个身份标识。删除消费者不影响组的last-delivered-id,也不影响其他消费者的进度。
三、销毁与重置:DESTROY和SETID的使用场景
当需要彻底废弃一个消费组时,用XGROUP DESTROY key group。这个操作会连带着清空该组的所有Pending消息记录,属于不可逆操作,执行前务必确认PEL里没有还没处理完的关键消息。可以先执行XPENDING key group确认挂起数量为0,再做销毁。
还有一个进阶命令XGROUP SETID key group id|$,它可以在不删除组的情况下重置消费位点。比如线上出了bug导致大量消息处理失败,修复后想重新消费某段历史消息,就可以用SETID把位点拨回去。注意这个操作同样不会清空PEL,旧位点之前已经pending的消息依然在列,需要配合XCLAIM或者干脆DESTROY重建来处理。
127.0.0.1:6379> XINFO GROUPS orders
1) 1) "name"
2) "order-group"
3) "consumers"
4) (integer) 2
5) "pending"
6) (integer) 0
7) "last-delivered-id"
8) "1718000000000-3"
127.0.0.1:6379> XGROUP SETID orders order-group 0
OKSETID执行后,组会回到从头消费的状态,配合XREADGROUP就能重新拉取历史消息做补偿处理。这在数据修复、故障回放的场景里非常实用。
四、消费者组的消费与确认机制
创建好组之后,消费和确认是日常操作的核心。XREADGROUP GROUP group consumer COUNT 10 BLOCK 5000 STREAMS key >中的>表示读取从未投递给任何消费者的新消息,而如果写具体的消息ID,则表示读取该消费者自己PEL中的消息,用于重试本地失败的任务,这两种模式一定要分清。
消费完成后必须执行XACK key group id把消息从PEL中移除,否则消息会一直处于挂起状态。完整的健康消费模型应该是:读取、处理、确认,再配合一个定时任务轮询XPENDING,对超时未确认的消息执行XAUTOCLAIM转移,形成闭环。
127.0.0.1:6379> XREADGROUP GROUP order-group worker-1 COUNT 10 STREAMS orders >
1) 1) "orders"
2) 1) 1) "1718000000001-0"
2) 1) "event"
2) "paid"
127.0.0.1:6379> XACK orders order-group 1718000000001-0
(integer) 1最后提一点容量规划:Stream本身可以用MAXLEN控制长度,但要注意消费者组的位点只推进不回退,如果消息被修剪掉而组还没消费到,这些消息就永久丢失了。所以设置MAXLEN时建议给消费延迟留足余量,或者改用近似裁剪XADD key MAXLEN ~ 10000 *降低性能开销。掌握XGROUP这一整套命令,再配合XACK和XPENDING,用Redis实现一套可靠的消息处理流程就完全够用了。