Redis 消息队列方案:从 List 到 Streams 的演进
文章目录

分布式系统里,消息队列是组件间异步通信的基础设施。生产者把消息丢进队列就可以继续干活,消费者按自己的节奏从队列里取消息处理。这种解耦让系统在流量洪峰时不至于雪崩。Redis 凭借高速读写的特性,天然适合做轻量级消息队列。但"适合"不等于"完美",消息队列有三个硬性需求,Redis 的不同数据类型对这三个需求的满足程度各不相同。
消息队列的三大核心需求

消息保序
消费者必须按照生产者发送消息的顺序来处理。以库存更新为例:生产者先发"库存设为 5",再发"库存设为 3"。如果消费者先处理了第二条,最终库存就错了。即使把消息改成增量形式(扣减 5、扣减 2),一旦中间夹杂读操作,乱序仍然会导致读到错误的中间状态。
重复消息处理
网络抖动可能导致消息重传,消费者收到两条一模一样的消息。如果不做幂等处理,同一条"扣减库存 5"被执行两次,库存就多扣了。解决方案是给每条消息分配全局唯一 ID,消费者维护已处理 ID 集合,收到重复 ID 直接跳过。
消息可靠性
消费者取出消息后还没处理完就宕机了,这条消息不能就此丢失。队列需要有某种机制让消费者重启后能重新拿到未完成的消息。
基于 List 的消息队列
List 是 Redis 最朴素的队列实现。LPUSH 从左端写入,RPOP 从右端取出,天然保证 FIFO 顺序。
# 生产者写入
LPUSH mq "101030001:stock:5"
# 消费者读取
RPOP mq
阻塞读取
如果消费者用 while 循环不停 RPOP,队列为空时会空转浪费 CPU。Redis 提供了 BRPOP 命令,队列为空时阻塞等待,有新消息写入时立即返回。这比轮询高效得多。
BRPOP mq 0
# 0 表示无限等待,直到有消息
全局唯一 ID
List 本身不会为消息生成 ID,需要生产者自行生成后拼入消息体。比如用雪花算法或者 UUID 生成 ID,再和业务数据拼成 JSON 或分隔符格式写入 List。消费者取出后解析 ID,对比已处理集合做幂等判断。
可靠性保障:BRPOPLPUSH
RPOP 取出消息后,消息就从 List 里消失了。如果消费者此时宕机,消息就丢了。BRPOPLPUSH 命令在取出消息的同时,把消息复制一份到另一个 List(备份队列):
BRPOPLPUSH mq mq_backup 0
消费者正常处理完后,从备份队列删除对应消息。如果宕机重启,从备份队列重新读取未处理的消息。代价是需要额外维护备份队列的清理逻辑。
List 方案的局限

List 最大的短板是不支持消费组。当生产者写入速度远超单个消费者的处理速度时,我们希望多个消费者分担消息处理。但 List 的 RPOP 是竞争式的——多个消费者同时 BRPOP 同一个 List,每条消息只会被其中一个消费者拿到。这看起来像消费组,但缺少消息确认机制、缺少消费进度追踪、缺少消息重投递能力。

另一个问题是不支持多消费组。如果同一条消息需要被多个不同的消费者组分别处理(比如一条数据既要做实时计算又要写入 HDFS),List 无法满足。
基于 Streams 的消息队列
Redis 5.0 引入的 Streams 是专门为消息队列设计的数据类型,借鉴了 Kafka 的设计理念。它在 List 的基础上补齐了消费组、ACK 机制、消息 ID 自动生成等关键能力。
写入:XADD
XADD mqstream * repo 5
"1599203861727-0"
* 让 Redis 自动生成全局唯一 ID,格式是 毫秒时间戳-序号。同一毫秒内的多条消息通过递增序号区分。也可以手动指定 ID,但自动生成更省心。
消息内容是键值对形式,一条消息可以包含多个字段:
XADD mqstream * device_id 33 temperature 26.8 pressure 101.3
读取:XREAD
# 从指定 ID 之后开始读取
XREAD BLOCK 100 STREAMS mqstream 1599203861727-0
# 读取最新消息,阻塞等待
XREAD BLOCK 10000 STREAMS mqstream $
BLOCK 参数实现阻塞读取,单位毫秒。$ 表示只读取调用之后新到达的消息。
消费组:XGROUP + XREADGROUP
消费组是 Streams 的核心差异化能力。创建消费组:
XGROUP CREATE mqstream group1 0
# 0 表示从头开始消费
消费组内的消费者用 XREADGROUP 读取:
XREADGROUP GROUP group1 consumer1 COUNT 1 STREAMS mqstream >
> 表示读取尚未被该消费组消费过的消息。关键规则:同一消费组内,一条消息只会被一个消费者读取。不同消费组之间互不影响,每个组都能独立消费全量消息。
这就解决了"多消费组"的需求:创建 group1 和 group2 两个消费组,group1 里的 consumer 做实时计算,group2 里的 consumer 写 HDFS,互不干扰。
同一消费组内多个消费者则实现了负载均衡:
XREADGROUP GROUP group2 consumer1 COUNT 1 STREAMS mqstream >
XREADGROUP GROUP group2 consumer2 COUNT 1 STREAMS mqstream >
XREADGROUP GROUP group2 consumer3 COUNT 1 STREAMS mqstream >
三个消费者各自拿到不同的消息,分担处理压力。
ACK 机制:XPENDING + XACK
Streams 为每个消费组维护一个 PENDING List,记录已被消费者读取但尚未确认的消息。消费者处理完消息后发送 XACK 确认:
XACK mqstream group2 1599274912765-0
如果消费者宕机重启,可以用 XPENDING 查看哪些消息还没确认:
# 查看 group2 的整体 pending 情况
XPENDING mqstream group2
# 查看 consumer2 的具体 pending 消息
XPENDING mqstream group2 - + 10 consumer2
未确认的消息不会丢失,消费者重启后可以重新处理。这比 List 的 BRPOPLPUSH 方案优雅得多,不需要手动维护备份队列。
List vs Streams 对照
| 能力 | List | Streams |
|---|---|---|
| 消息保序 | 支持(FIFO) | 支持(FIFO) |
| 全局唯一 ID | 需要生产者自行生成 | 自动生成 |
| 阻塞读取 | BRPOP | XREAD BLOCK |
| 消费组 | 不支持 | 原生支持 |
| 多消费组 | 不支持 | 支持 |
| ACK 机制 | 需要 BRPOPLPUSH + 手动管理 | XPENDING + XACK |
| 消息回溯 | 不支持(读即删) | 支持(消息持久保留) |
| 版本要求 | 所有版本 | Redis 5.0+ |
Redis 做消息队列的边界
Redis 做消息队列的优势是轻量、快速、部署简单。一个 Redis 进程就是一个消息队列,不需要像 Kafka 那样还得部署 ZooKeeper(虽然 Kafka 新版本已经在去 ZK)。对于消息通信量不大、对可靠性要求不是极致的场景,Redis 是非常好的选择。
但 Redis 做消息队列有几个绕不开的短板:
数据可靠性。Redis 的持久化是异步的。AOF 配置为 everysec 时,宕机最多丢 1 秒数据。主从切换时如果从库同步有延迟,也可能丢消息。这对金融交易类场景是不可接受的。
消息堆积能力。Redis 是内存数据库,消息堆积意味着内存持续增长。如果消费者长时间跟不上生产者,内存会被撑爆。Kafka 的消息存在磁盘上,堆积 TB 级数据也不影响性能。
功能完备性。延迟队列、死信队列、消息优先级、事务消息这些高级特性,Redis 原生都不支持,需要业务层自己实现。
所以选型建议是:
- 发短信、发通知、日志采集这类对丢消息不敏感的场景,Redis 完全够用,轻量高效。
- 交易、支付、订单这类对数据完整性要求极高的场景,用 Kafka 或 RabbitMQ。
- 消息量极大、需要长期堆积的场景,用 Kafka。
如果你的 Redis 版本是 5.0 以上,优先用 Streams。如果是老版本且不方便升级,List + BRPOPLPUSH 也能凑合用,但要做好幂等和备份队列的维护工作。
实践中的几个注意点
第一,消费者的幂等设计。不管用 List 还是 Streams,消费者都必须做幂等。网络重传、消费者重启后重新消费 PENDING 消息,都会导致同一条消息被处理多次。幂等的实现方式取决于业务:数据库写入可以用唯一索引,状态变更可以用版本号,计数类可以用 Redis 的 SETNX 做去重标记。
第二,Streams 的内存管理。Streams 的消息默认不会自动删除(即使所有消费组都 ACK 了)。需要用 MAXLEN 参数限制 Stream 长度,或者定期用 XTRIM 裁剪:
# 写入时限制最大长度
XADD mqstream MAXLEN ~ 10000 * repo 5
# 手动裁剪
XTRIM mqstream MAXLEN 10000
~ 是近似裁剪,性能比精确裁剪好。
第三,消费者故障检测。XPENDING 可以看到每条 pending 消息的空闲时间(idle time)。如果某个消费者的消息 idle 时间过长,说明它可能已经挂了。可以用 XCLAIM 命令把这些消息转移给其他消费者处理:
XCLAIM mqstream group1 consumer2 3600000 1599274912765-0
# 把 idle 超过 1 小时的消息转给 consumer2
第四,监控指标。生产环境里要监控 Stream 的长度(XLEN)、各消费组的 lag(未消费消息数)、pending 消息数。lag 持续增长说明消费者跟不上,需要扩容消费者或者排查消费逻辑的性能瓶颈。
更多推荐


所有评论(0)