3.3 Stream:可持久化的消息日志 本节摘要:Stream 是 append-only 的日志结构消息类型,底层基于 RADIX 树按消息 id 索引。相比 List 队列,它补上了消费组、待确认列表与消息重放三大能力,把 Redis 队列从"能凑合"升级到"可认真用"。 List 队列缺什么 2.2 留下的问题:BRPOP 弹出即删除,消费者处理到一半崩溃,消息就没了;两个消费者抢同一个 List,消息被随机分掉,无法各自独立消费。这两个缺口分别是确认机制与消费组,Stream 的结构就是冲着它们设计的。 消息写入只追加在日志尾部,每条消息的 id 形如"毫秒时间戳减序号",天然有序且可寻址。
本节摘要:Stream 是 append-only 的日志结构消息类型,底层基于 RADIX 树按消息 id 索引。相比 List 队列,它补上了消费组、待确认列表与消息重放三大能力,把 Redis 队列从"能凑合"升级到"可认真用"。
2.2 留下的问题:BRPOP 弹出即删除,消费者处理到一半崩溃,消息就没了;两个消费者抢同一个 List,消息被随机分掉,无法各自独立消费。这两个缺口分别是确认机制与消费组,Stream 的结构就是冲着它们设计的。
消息写入只追加在日志尾部,每条消息的 id 形如"毫秒时间戳减序号",天然有序且可寻址。消费组只是三个数据结构的组合:记录组内已投递位置的游标、待确认消息的 id 列表、以及消费者名单。
这三个部件的状态全部挂在 Stream 键本身上,删键等于连组带游标一起清空;日常运维想只清消费进度不动日志,用 XGROUP SETID 把游标拨回起点,或者 XGROUP DESTROY 删组重建——理解"组不是独立的键",误删与误建的行为就都能预判。

# 生产:XADD 追加,* 表示让服务端生成 id > XADD orders * type create uid 1001 amt 299 "1724290000000-1" # 建组:从日志最早开始消费 > XGROUP CREATE orders g1 0 # 消费:组内取一条,超时 5 秒阻塞等待 > XREADGROUP GROUP g1 c1 COUNT 1 BLOCK 5000 STREAMS orders > 1) 1) "orders" 2) 1) 1) "1724290000000-1" 2) 1) "type" 2) "create" # 处理成功,确认 > XACK orders g1 1724290000000-1 (integer) 1
建组命令还有个常用变体:XGROUP CREATE orders g1 $ 以日志当前末尾为起点,只消费建组之后的新消息——回填历史与只听新增,差别就在这一个参数。组建成后,XINFO GROUPS orders 一眼看完所有组的运行状态:每个组的消费者数、积压数、游标位置,巡检脚本扫这张表就够。
关键是 > 与 0 两个符号:> 表示"只给我游标之后的新消息",0 表示"把我待确认列表里没 ACK 的重新给我"。崩溃恢复就是用 0 重读——未确认消息不会消失,它们躺在 PEL 里等着。
id 的规则展开说:毫秒时间戳加短横线加序号。同一毫秒内多条消息靠序号递增区分,所以 id 全局严格递增、天然可比较。两个实用变式:XADD 时可以自己指定完整 id(比如从业务流水号映射),服务端会校验必须大于当前最后一条;XREAD 支持按任意 id 定位读取,XREAD STREAMS orders 0 从头读、传一个具体 id 则从它之后读——"消息日志可重放"的入口就是这两个参数,事件溯源、对账补数、新下游上线回填历史,全靠它。
消费者挂掉后,PEL 里的消息需要有人接管:
# 把 c1 挂起超过 60 秒的消息转移给 c2,重新投递 > XAUTOCLAIM orders g1 c2 60000 0 # 查看组内积压情况 > XPENDING orders g1 # 组内消息堆积超限则修剪日志,保留最近 100 万条 > XTRIM orders MAXLEN 1000000
XPENDING 的完整输出值得逐段读:
> XPENDING orders g1 1) (integer) 3 # 组内未确认总数 2) "1724290000000-1" # 最早的未确认id 3) "1724290005000-2" # 最晚的未确认id 4) 1) 1) "c1" # 各消费者的积压数 2) (integer) 3
四行的用法:第一行是积压总量的监控指标,接告警;第二三行圈出积压的时间范围,判断是"卡在很久前"还是"刚刚的";第四行定位到具体消费者。监控脚本拿这张表做两件事——总积压超阈值告警、某消费者的条目长时间不动就触发 XAUTOCLAIM 接管。日志不能无限涨,XTRIM 按条数或最小 id 裁剪,这是运维上必须定时的动作。
Stream 补齐了 Redis 队列的语义短板,但与 Kafka、RabbitMQ 相比仍有硬边界:内存介质决定容量有限(消息要靠 XTRIM 约束);主从复制是异步的,故障切换瞬间可能丢少量已确认消息(第 7 章会看到复制原理);没有完备的跨机房复制故事。我的实践判断:日百万级以内、允许极端情况下少量重投的消息流,Stream 完全胜任;金融级不丢不重,交给专业 MQ。
消费者怎么横向扩容,值得单独说清。一个消费组内的消息只会投给组内某个消费者——组是"竞争消费",不同组才是"独立消费"。所以扩容的正确姿势是给同一组加消费者实例(比如 worker-1 到 worker-8),消息自动在它们之间分摊;XINFO CONSUMERS 能看到每个消费者分到了多少、积压多少,负载不均时先查这一步。减容时把退役消费者的 PEL 用 XAUTOCLAIM 转移给存活者,再删消费者,积压才不会跟着尸体一起失踪。
另一个必谈的边界是投递语义:Stream 的组合机制给出的是"至少一次"——消息确认前不删,接管会重投,崩溃恢复也会重读。这意味着下游必须幂等:以消息 id 或业务单号做去重键,处理前先查去重表(一个带 TTL 的 SET NX 就够)。指望"恰好一次"而省掉幂等层,重复消息迟早在线上教做人。
和专业消息队列的差距还体现在"背压"上:Kafka 的消费组滞后可以反馈到生产端感知;Stream 的 XADD 永远成功(只要内存没满),生产端对下游堆积毫无感知,堆积全靠消费侧自己用 XPENDING 与 XINFO 盯。把这两个命令的输出接进告警,是 Stream 方案上线前的最后一道必做工序。
⚠️ 常见坑:消费者忘了 XACK 是最常见事故,PEL 无限膨胀直到内存告警。务必给所有消费路径配异常分支的 ACK 或转移逻辑,并用 XPENDING 监控积压时长。