Redis Stream从5.0版本引入后,逐渐成为轻量级消息队列的热门选择。它支持持久化、支持按ID范围读取,但单靠XREAD命令读取有个明显短板:所有客户端读到的是同一份全量数据,无法在多个实例之间做负载均衡。XREADGROUP解决了这个问题,它配合消费者组机制,让同一个Stream的消息自动分发给组内不同的消费者,实现真正的队列语义。

一、消费者组的基本原理
消费者组可以理解为Stream之上的一个逻辑视图。创建消费者组时,Redis会在内部记录三个关键信息:该组的最后投递ID(last-delivered-id)、组内每个消费者的pending列表、以及组名本身。当消费者通过XREADGROUP读取时,Redis只会把last-delivered-id之后的消息投递给它,并自动更新这个游标。
假设有消费者A和B同属一个组,A读取时拿到消息1、2、3,B再读取时拿到的就是4、5、6,两组数据互不重叠。这就是组内消息分发的核心:同一条消息在一个组内只会被投递给一个消费者。而不同的消费者组之间互不影响,如果创建两个组group-a和group-b,两边各自都能读到全量消息,常用于一个事件被多个业务方订阅的场景。
pending列表是消费者组的另一大设计亮点。消息投递给消费者后不会立即删除,而是进入该消费者的pending列表,直到消费者显式调用XACK确认。如果消费者处理到一半崩溃了,重启后可以继续从pending列表读取未确认的消息,保证消息不丢失。这种设计本质上借鉴了传统MQ的ack机制,用较小的存储代价换来了可靠消费能力。
二、XREADGROUP命令完整用法
先看命令的标准格式:XREADGROUP GROUP group consumer [COUNT n] [BLOCK ms] [NOACK] STREAMS key [key ...] id [id ...]。GROUP后面跟组名和消费者名,消费者不需要预先创建,第一次使用时自动注册。id参数一般传特殊符号>,表示只读取该组从未投递过的新消息。
下面用redis-cli演示一个完整流程:
$ # 创建Stream并写入几条测试消息
$ XADD mystream * task send-email
"1690000000000-0"
$ XADD mystream * task send-sms
"1690000000000-1"
$ # 创建消费者组,从头开始消费
$ XGROUP CREATE mystream mygroup 0
OK
$ # 消费者worker-1读取新消息
$ XREADGROUP GROUP mygroup worker-1 COUNT 10 STREAMS mystream >
1) 1) "mystream"
2) 1) 1) "1690000000000-0"
2) 1) "task" 2) "send-email"
2) 1) "1690000000000-1"
2) 1) "task" 2) "send-sms"
$ # 确认第一条消息处理完成
$ XACK mystream mygroup 1690000000000-0
(integer) 1如果把id位置的>换成0或者具体ID,含义就变了:此时读取的是当前消费者自己pending列表中ID大于给定值的消息,也就是已经被投递但未确认的历史消息。这个特性常用于消费者重启后恢复现场,代码里通常会先扫描pending列表处理残留消息,再用>读新消息。
BLOCK参数用于阻塞式读取,比如BLOCK 5000表示没有新消息时等待5秒,超时返回nil。写常驻消费进程时一般配合BLOCK 0(无限阻塞)使用,配合socket超时设置,既省去轮询开销又能保活连接。NOACK参数则表示读取即确认,消息不进pending列表,适合允许少量丢失、追求吞吐的场景,比如日志采集。
三、pending消息与故障恢复
消费者处理消息途中宕机是常态,pending机制就是为了这种情况准备的。用XPENDING可以查看组内所有未确认消息的概览:
$ XPENDING mystream mygroup
1) (integer) 1 # pending总数
2) "1690000000000-1" # 最小ID
3) "1690000000000-1" # 最大ID
4) 1) 1) "worker-1" # 消费者及各自的pending数量
2) (integer) 1
$ # 查看详细信息,包括投递次数和空闲时间
$ XPENDING mystream mygroup - + 10
1) 1) "1690000000000-1"
2) "worker-1"
3) (integer) (integer) 3600000 # 空闲毫秒数
4) (integer) 2 # 已投递次数发现某条消息长时间没被确认,就要考虑转移所有权。XCLAIM命令可以把pending消息从故障消费者转移到健康消费者:
$ # 把空闲超过60秒的消息转移给worker-2,只转移一条 $ XCLAIM mystream mygroup worker-2 60000 1690000000000-1 1) 1) "1690000000000-1" 2) 1) "task" 2) "send-sms"
空闲时间参数很关键,只有超过该时长的消息才会被转移,避免误抢正在处理中的消息。Redis 6.2之后还提供了XAUTOCLAIM命令,它相当于XPENDING加XCLAIM的组合,自动扫描并认领符合条件的消息,用起来更省事。实际编码中建议每个消费者循环里定期执行认领逻辑,把超时消息拉回处理队列。
还有一个细节值得注意:投递次数可以用来判断毒消息。如果某条消息被投递了很多次仍未确认,大概率是处理逻辑报错导致的,此时应该把它转入死信处理,而不是无限重试。可以在业务代码里对delivery count设置阈值,超过后写入另一个专门记录异常的Stream。
四、生产环境的使用建议
首先是消费位移的初始化问题。XGROUP CREATE的最后一个参数决定从哪里开始读:$表示只消费创建之后的新消息,0表示从头消费全部历史消息。如果组已存在会报BUSYGROUP错误,代码里要做好异常捕获,很多客户端封装了MKSTREAM选项,可以在Stream不存在时自动创建,省去初始化顺序的烦恼。
其次是Stream的清理策略。Stream不像List消费完就删,必须显式删除。常用的方案是定期执行XTRIM mystream MAXLEN 10000限制长度,或者用MINID按时间裁剪。要注意XTRIM是按ID删除,不看消费者组的进度,如果裁剪太激进,落后的组就再也读不到被删的消息了,所以裁剪阈值要大于业务允许的最大消费延迟。
最后是幂等性的设计。消费者组保证了at-least-once语义,消息至少被投递一次,但崩溃恢复、消息认领都可能导致同一条消息被处理多次。业务侧必须做幂等,常用手段是在消息里带唯一业务ID,处理前先查Redis或数据库的去重标记。不要指望消息队列本身保证恰好一次投递,把幂等做在消费端才是稳妥的做法。
总的来说,XREADGROUP加消费者组是Redis Stream最实用的能力,掌握pending、ack、claim这套机制后,用Redis搭建轻量级、高可用的异步任务系统并不复杂。当业务量级还没到必须引入Kafka或RabbitMQ时,这套方案能显著降低运维成本。
RedisXREADGROUP消费者组修改时间:2026-09-11 08:54:35