Redis Stream是Redis 5.0引入的一种追加式日志数据结构,专门用来实现轻量级消息队列。它把每条消息作为日志项追加到流中,每条消息拥有全局唯一且递增的ID,消费者可以通过游标读取指定区间的数据。相比早期的列表或发布订阅方案,Stream在消息持久化、多消费者协同、未确认消息追踪上提供了原生支持,非常适合中小型业务做异步任务解耦。

Stream底层结构与基础写入读取原理
Stream在Redis内部使用基数树(Rax树)来组织消息ID与内容的映射,消息ID默认由毫秒时间戳和序列号组成,例如1715000000000-0。这种结构让范围查询非常高效,消费者可以按ID区间快速定位。写入时使用XADD命令,如果省略ID则Redis自动生成,保证严格递增且不会回退,从而避免消息乱序。
单纯写入并不能构成队列,还需要读取端配合。最基本的读取命令是XREAD,它可以阻塞等待新消息到达。但XREAD只是单消费者视角,无法记录“谁读了哪条”,因此生产环境通常使用消费者组。消费者组通过XGROUP CREATE创建,之后用XREADGROUP读取,Redis会为组内每个消费者维护一个待确认列表(PEL),只有执行XACK后才算真正消费完成。
下面是一段基础的写入与读取示例,展示了自动生成ID和阻塞读取的用法:
# 向名为 mystream 的流写入一条消息,字段为 msg,值为 hello XADD mystream * msg hello # 从0号ID开始非阻塞读取最多10条 XREAD COUNT 10 STREAMS mystream 0 # 阻塞等待新消息,使用 $ 表示只接收最新之后的消息 XREAD BLOCK 5000 COUNT 10 STREAMS mystream $
消费者组机制与可靠消费实现
消费者组是Redis Stream解决“至少一次投递”的核心。当多个消费者加入同一个组,Redis会把流中的消息轮流分发给不同消费者,并记录在PEL里。如果某个消费者崩溃,它名下未确认的消息不会丢失,其他消费者可以通过XCLAIM命令认领这些pending消息继续处理,从而保证任务不中断。
创建消费者组时需要指定起始ID,通常用0表示从开头消费,或用$表示只消费新消息。读取时必须携带组名和消费者名,例如XREADGROUP GROUP group1 consumer1 COUNT 5 STREAMS mystream >,其中>表示读取尚未分配给本消费者的消息。处理成功后务必调用XACK,否则PEL会无限增长,导致内存泄漏和重启后重复消费。
以下代码演示了完整的组消费与确认流程,包含异常情况下查看pending消息的方式:
# 创建消费者组,从开头消费 XGROUP CREATE mystream group1 0 # 消费者A读取消息 XREADGROUP GROUP group1 A COUNT 2 STREAMS mystream > # 假设读到 1715000000000-0,处理完后确认 XACK mystream group1 1715000000000-0 # 查看组中所有 pending 消息 XPENDING mystream group1 # 将超时未确认的消息转移给消费者B XCLAIM mystream group1 B 30000 1715000000001-0
在实际部署中,建议给XCLAIM设置合理的空闲时间(如30秒),避免消息刚在网络延迟中被重复认领。同时业务处理要保持幂等,因为网络分区可能导致同一条消息被投递两次。
性能调优与常见误区分析
很多团队把Stream当作无限增长的日志,长期不清理导致内存暴涨。Redis提供了XTRIM命令来限制流长度,例如XTRIM mystream MAXLEN 100000可以近似保留最近十万条。对于已完成且无需回溯的消息,主动修剪既能降低内存,也能加快范围扫描速度。
另一个误区是忽视PEL堆积。如果消费者频繁崩溃却不修复pending,XPENDING输出会越来越大,影响XREADGROUP的响应。应当编写监控脚本,定期将空闲超过阈值的消息重新分配,并在业务层做去重表(如Redis的SET记录已处理ID)来防止重复执行。
在超高并发场景下,单个Stream的写入可能成为瓶颈。此时可以按业务键做流分片,比如把用户ID取模分散到mystream_0、mystream_1等多个流,再配合多个消费者组提升并行度。下面给出一段Java风格的逻辑片段,说明如何根据分片读写:
// 计算分片索引
int shard = userId % 4;
String streamKey = "mystream_" + shard;
// 写入消息
jedis.xadd(streamKey, StreamEntryID.NEW_ENTRY, Map.of("data", payload));
// 对应消费者组读取
jedis.xreadGroup(group, consumer,
Map.of(streamKey, StreamEntryID.UNRECEIVED_ENTRY),
100, 0);
最后需要注意,Stream虽然支持阻塞读,但不具备严格的全局顺序保证跨分片事务。如果业务要求强一致,应在上层用数据库事务或幂等补偿机制兜底,而非完全依赖消息队列本身。
Redis_Stream消息队列消费者组修改时间:2026-08-13 21:30:28