Redis 的 Stream 数据结构和消息队列应用

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 还有一个容易被忽略的能力:上限控制。通过 MAXLENMINID 可以裁剪历史消息,避免内存无限增长:

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 可靠性的关键:消费者崩溃后,消息不会丢失,其他消费者可以通过 XCLAIMXAUTOCLAIM 接管。

查看待处理消息

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 裁剪,或者监控 XLENXPENDING。积压过多会撑爆内存。

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 数据结构和消息队列应用

赞 (0) 打赏

评论 0

取消
  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址

觉得文章有用就打赏一下文章作者

支付宝扫一扫打赏

微信扫一扫打赏