Redis 在 5.0 版本引入 Stream 类型,目标很明确:提供一个更完整、更可靠的消息队列原语。在此之前,开发者用 List 做队列、用 Pub/Sub 做广播,但都有明显短板。Stream 填补了持久化、消费者组、消息确认这些关键能力的空白。下面从面试常见问题出发,把 Stream 的结构、核心命令和消息队列实践串起来。
为什么需要 Stream
先说清楚它解决了什么问题。
List 做队列的痛点:LPUSH + BRPOP 可以实现简单的 FIFO 队列,但消息一旦被取出就从 List 中删除了。如果消费者处理失败,消息就丢了。没有 ACK 机制,没有消费者组,多个消费者之间也无法方便地分摊消息。
Pub/Sub 的痛点:消息是即发即弃的,没有持久化。消费者离线期间的消息直接丢失,也不支持回溯消费。
Stream 的定位:它像一个 append-only 的日志,每条消息有唯一 ID,支持持久化、消费者组、ACK 确认、消息回溯。可以理解为 Redis 内部实现了一个简化版的 Kafka。
Stream 的核心结构
Stream 本质上是一个有序的消息链表,每条消息包含:
- 消息 ID:格式为
时间戳-序列号,如1686000000000-0。ID 单调递增,保证顺序。 - 消息内容:一个 field-value 对的集合,类似 Hash 结构。
XADD mystream * name Alice age 30
这里 * 表示让 Redis 自动生成 ID。返回值就是这条消息的 ID。
Stream 还有一个容易被忽略的能力:上限控制。通过 MAXLEN 或 MINID 可以裁剪历史消息,避免内存无限增长:
XADD mystream MAXLEN ~ 1000 * name Bob
~ 表示近似裁剪,性能更好,不会精确到每一条。
消费者组:Stream 做消息队列的核心
单独用 Stream 只是读日志,真正让它成为消息队列的是消费者组(Consumer Group)。
创建消费者组
XGROUP CREATE mystream mygroup $ MKSTREAM
$表示从最新消息开始消费,0表示从头开始。MKSTREAM表示 Stream 不存在时自动创建。
消费者读取消息
XREADGROUP GROUP mygroup consumer1 COUNT 10 STREAMS mystream >
> 表示读取从未投递给该组的新消息。Redis 会把消息标记为「已投递给 consumer1」,但不会删除。
ACK 确认
消费者处理完成后,必须显式确认:
XACK mystream mygroup 1686000000000-0
未被 ACK 的消息会留在 PEL(Pending Entries List,待处理条目列表) 中。这是 Stream 可靠性的关键:消费者崩溃后,消息不会丢失,其他消费者可以通过 XCLAIM 或 XAUTOCLAIM 接管。
查看待处理消息
XPENDING mystream mygroup
可以查看有多少消息未被确认、分别属于哪个消费者、空闲了多久。面试中经常问:如何发现消费积压? 答案就是通过 XPENDING 监控 PEL 的长度和消息的空闲时间。
消息回溯与重读
Stream 支持按 ID 范围读取历史消息:
XRANGE mystream - +
XREAD COUNT 100 STREAMS mystream 0
- 和 + 分别表示最小和最大 ID。这个能力在排查问题、重新处理数据时非常有用,也是 List 和 Pub/Sub 做不到的。
Stream vs Kafka vs RabbitMQ
面试中常被拿来对比。核心差异:
| 维度 | Redis Stream | Kafka | RabbitMQ |
|---|---|---|---|
| 持久化 | 内存 + RDB/AOF | 磁盘 | 磁盘 |
| 吞吐量 | 高(受内存限制) | 极高 | 中高 |
| 消息保留 | 可配置上限 | 按时间/大小保留 | 消费后删除 |
| 消费者组 | 支持 | 支持 | 支持 |
| 运维复杂度 | 低 | 高 | 中 |
| 适用场景 | 轻量级队列、实时性要求高 | 大数据管道、日志 | 企业级消息中间件 |
结论:Stream 适合中小规模、对延迟敏感、不想引入额外中间件的场景。如果消息量极大或需要长期存储,Kafka 更合适。
实战中的关键注意事项
1. 消费者组不能自动创建消费者
消费者需要显式用 XREADGROUP 注册。如果消费者长期离线,它的 PEL 消息需要被其他消费者认领。
2. 避免消息无限积压
设置 MAXLEN 裁剪,或者监控 XLEN 和 XPENDING。积压过多会撑爆内存。
3. ACK 的时机
一定要在业务处理成功后再 ACK。先 ACK 再处理,崩溃就丢消息;处理完再 ACK,最多重复消费,配合幂等设计即可。
4. 消费者命名要稳定
用固定的消费者名称(如主机名 + 进程 ID),否则每次重启都产生新消费者,PEL 会越来越乱。
5. XAUTOCLAIM 是 Redis 6.2+ 的改进
旧版用 XCLAIM 逐条认领,效率低。XAUTOCLAIM 可以批量认领超时消息,更适合故障恢复。
一个典型的消费循环
while True:
# 读取新消息
messages = redis.xreadgroup(
'mygroup', 'consumer1',
{'mystream': '>'},
count=10, block=5000
)
for stream, msgs in messages:
for msg_id, data in msgs:
try:
process(data)
redis.xack('mystream', 'mygroup', msg_id)
except Exception:
# 不 ACK,留给重试或认领
pass
这个模式的关键点:处理成功才 ACK,失败的消息留在 PEL 中等待后续处理。
总结
Redis Stream 的核心价值在于:在 Redis 内部提供了一个具备持久化、消费者组、ACK 确认和消息回溯能力的消息队列。它不适合替代 Kafka,但在轻量级场景下,省去了引入独立消息中间件的成本。面试中重点掌握:消息 ID 结构、消费者组的工作机制、PEL 与 ACK 的关系、以及和 List/Pub/Sub 的本质区别。理解这些,基本能覆盖 Stream 相关的绝大多数问题。
未经允许不得转载:任鹏个人博客 » Redis 的 Stream 数据结构和消息队列应用

