尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Redis Stream消息队列深入解析:原理、消费者组与可靠性实践

Redis Stream消息队列深入解析:原理、消费者组与可靠性实践 前两天一个读者跑来跟我说他面一家中厂的时候被问到“Redis 5.0 的 Stream 消息队列了解吗”他大概说了几句“Stream 是 Redis 5.0 新增的消息队列、支持消费者组、消息不会丢”然后就卡壳了。面试官接着追问“那它和 Pub/Sub 有什么区别消费者组是怎么工作的消息 ID 是怎么生成的”他就答不上来了。这个场景太典型了。Stream 不是说背几个命令就能糊弄过去的面试官真正想听的是你知不知道 Stream 解决什么问题、底层结构长什么样、消费者组的消息确认机制怎么运作、跟 Kafka 这类消息中间件比有什么优势和劣势。说白了他想确认你是真的在项目里用过、踩过坑还是只在八股文里见过。这篇文章我想从原理、命令、可靠性和线上场景几个角度把 Stream 说透尽量用我在实际项目里摸爬滚打的经验来讲。不管你是准备面试还是正在为项目做技术选型看完应该能对 Stream 有个完整、落地的认知。1. 面试官在问 Stream 时到底想听到什么1.1 Stream 是为解决什么问题才被设计出来的在 Redis 5.0 之前如果你想用 Redis 做消息队列无非两条路List 和 Pub/Sub。List 方案就是LPUSHBRPOP实现一个简单的阻塞队列。它能存数据、能阻塞读、能做粗略的负载均衡但硬伤很明显不支持多消费者组。多个消费者同时 BRPOP 同一个列表一条消息只会被一个消费者取走这叫做“竞争消费”类似一个线程池在抢任务。可如果你想要“一条消息被多个业务各自消费一次”比如订单服务要消费、风控服务也要消费List 就无能为力了。你只能开多个不同的 key 重复推那消息一致性就变成了你的噩梦。Pub/Sub 方案呢支持真正的发布订阅一条消息能广播给所有订阅者。但它是“发后即焚”的消息发出去如果没有订阅者在监听直接就丢了。消费者只要断线重连中间这阵子的消息就再也看不到了。对于日志、通知这类允许丢一下的场景还行但要拿来做业务消息队列一丢消息业务就要报警。Stream 出来之后这两个问题都补上了。它本质上是一个持久化的、支持多消费者组、支持消息确认和回溯的追加式日志结构。消息存在 Redis 内存里RDB 和 AOF 都能把它落盘消费者离线之后再回来还能从指定位置继续读。你可以开多个消费者组每个组之间互不干扰组内又可以挂多个消费者分摊组内的消息。1.2 面试官爱听的 Stream 和 Kafka 差异点不少人看到 Stream 的消费者组模型第一反应就是“这不就是简化版 Kafka 吗”。理解是对的但如果你在面试里这么答一定要把下面这层差异讲出来不然显得你是背过概念但没想明白。Kafka 的核心抽象是 partition一个消费者组内每个 partition 同时只能被一个消费者实例持有通过分区实现并发度和消息有序性。Redis Stream 里没有物理分区这个概念它只有一条追加日志。所谓的“消费者组”其实是在一条日志上维护了多个游标每个组有一个last_delivered_id组内的多个消费者通过竞争去读取新消息。组内多个消费者之间怎么分摊消息呢是靠命令执行时的抢占谁的XREADGROUP先到新消息就给谁而不是像 Kafka 那样把某个范围的数据绑定给某个消费者。所以 Stream 在单个 Stream key 上的并发伸缩能力天然不如 Kafka 的多分区模型。Redis 官方也给了方案如果业务量真的需要分区可以自己建多个 Stream key然后按业务维度做 hash 路由。这个后面讲实操的时候我会细说。面试时如果你能把“Kafka 分区逻辑是从存储层面隔离的Stream 的分区只是应用层自己路由出来的”这个观点讲出来面试官大概率会点个头。2. Stream 的底层结构与消息 ID 设计2.1 追加日志在内存里是怎么组织的你可以在 Redis 里简单地理解成 Stream 是一条只能追加的日志但底层不是一条链表而是一个叫做 Radix Tree基数树的结构树里的叶子节点用listpack紧凑编码来存消息。为什么不用简单链表因为链表做范围查询太痛苦了。Stream 最常见的操作是按消息 ID 区间查数据比如XRANGE链表为了找到中间某个位置只能从头遍历复杂度 O(n)。Radix Tree 是前缀树按 key 的公共前缀来压缩路径在内存里既能高效地做范围查询又对 Redis 的内存友好。消息 ID 是单调递增的在 Radix Tree 上追加新消息很快经典数据结构选型。listpack 是 Redis 用来替代 ziplist 的一种紧凑列表编码把多个连续的消息打包存到一个节点里每个消息只是其中的一个 entry。这样做既减少了指针数量也提高了内存利用率。对同一个 Stream 来说消息越多listpack 节点会越多Redis 会自动管理这些节点的创建和合并。这个底层机制不用背得很深但你要能说出来Stream 不是简单队列它是一棵有序的树支持按 ID 区间高效遍历底层用 listpack 压缩存储消息条目。这就能解释为什么 Stream 能在 Redis 弱内存的情况下还能承受大量消息的写入。2.2 消息 ID 的两段式设计到底好在哪Stream 每条消息都有一个全局唯一的 ID格式是时间戳-序号比如1716888888888-0。时间戳是 Redis 服务器本地毫秒时间序号是同一毫秒内的自增序号。两条消息不会生成同样的 ID即便时间戳回拨序号也会保证在同一毫秒内继续递增。这个 ID 设计有一个非常巧妙的点它把消息的时间属性和顺序属性合二为一。消费端拿到一条消息看到1716888888888-0就知道这条消息大概是什么时候写入的不用额外再存一个时间字段。范围查询上XRANGE mystream - 取全部XRANGE mystream 1716888888888 取某个时间点之后的消息做时间窗口统计或者故障恢复都特别方便。在面试的时候你还可以补一句Redis 也允许你自定义 ID比如XADD mystream 1716888888888-1 field value只要你给的 ID 比当前最大 ID 大就行。所以理论上你可以把业务系统里已有的全局 ID 当作消息 ID 塞进 Stream 里这样一来消费端可以直接用业务 ID 做去重不用维护 ID 映射关系。有些同学实际项目里就是这么玩的效果很好。2.3 消费者组在内部记录了什么创建消费者组时Redis 会为这个组维护两个核心的东西组级别的last_delivered_id以及每个消费者自己的 PELPending Entries List待确认消息列表。last_delivered_id代表这个组已经投递到哪儿了。消费者组里的消费者通过XREADGROUP ... 读取新消息时Redis 会原子性地更新这个游标。这样即使某个消费者读完还没处理完就挂了重启后组里其他消费者也能从last_delivered_id继续读新消息不会因为某个消费者掉线就把整条消费链路卡住。PEL 则是用来追踪“投递了但还没确认”的消息。当一个消费者读取了一条消息这条消息的 ID 就会进入该消费者的 PEL等它处理完执行XACKRedis 才把这条消息从 PEL 里移除。PEL 是 Stream 可靠投递的核心后面专门展开讲。你只需要记住一个关键结论Stream 的消费者的进度不是靠消费者自己记录消费到哪了而是靠 Redis 服务端的内部游标和确认状态来管理。这正是它比 List 复杂、也比 List 可靠的地方。3. 核心命令实操从生产到消费到消费组3.1 五条生产与读取命令一条条看透先看生产者侧。# 往有序流里追加消息 XADD order-events * event.created order_id 1001 amount 299.00 1716888888888-0*表示让 Redis 自动生成消息 ID返回值就是这条消息的 ID。后面跟的是 field-value 对你可以塞多个字段类似一个 hash。Stream 的每条消息本质就是一个小的 field-value map这一点和 Kafka 里只有 value blob 不一样。# 看长度 XLEN order-events (integer) 1 # 按区间读 XRANGE order-events - COUNT 10 # 反向读适合拿最新消息 XREVRANGE order-events - COUNT 10XRANGE的-和是最小 ID 和最大 ID 的简写配合 COUNT 可以一页一页翻。这种“按 ID 区间扫描”的能力是 List 和 Pub/Sub 完全不具备的做补偿任务、重放交易流水都用得上。再看消费者侧最常用的XREAD# 阻塞读新消息0 表示从最早开始$ 表示只读之后的 XREAD COUNT 1 BLOCK 5000 STREAMS order-events $BLOCK 5000是说没有新消息时最多阻塞 5 秒而不是无限阻塞STREAMS后面先写 key 名再写起始 ID。这里的$表示“只读我发起阻塞之后到达的新消息”。如果传入具体的消息 ID就会从该 ID 的下一条开始读这个特性可以用来做断点续读。3.2 消费者组创建与消费的核心套路消费者组才是 Stream 最精华的部分直接上实操。# 创建消费者组从最早消息开始消费 XGROUP CREATE order-events group_a 0 # 组内读消息consumer-1 是消费者名 XREADGROUP GROUP group_a consumer-1 COUNT 1 STREAMS order-events 这里是一个特殊 ID表示“给我这个组还没投递过的新消息”。如果没有新消息命令会返回空配合BLOCK参数就可以实现阻塞等待。组内读和普通读最大的区别消息一旦被某个组内消费者通过读走它会记录到该消费者的 PEL 里同组其他消费者就不会再读到同一条消息了。这就是消费者组内部的竞争消费模型。如果你传的不是而是一个具体消息 ID那就不是读新消息而是从自己的 PEL 里读取历史未确认消息。这个细节非常有用因为消费者崩溃重启后最先要处理的就是自己 PEL 里那些没确认的消息。# 消费完确认 XACK order-events group_a 1716888888888-0 (integer) 1XACK的第二个参数是组名第三个参数是消息 ID可以传多个。确认成功后消息就会从该消费者的 PEL 里移除。注意这里移除的只是“这个消费者组”的 PEL 记录并不影响 Stream 原始消息的存在。Stream 本身是一条长日志XACK只作用于消费组的进度状态不删除日志。3.3 查看积压情况XPENDING 和 XINFO线上排查消息积压这两个命令能救命。 XPENDING order-events group_a返回结果包含这个组总共有多少待确认消息、最早和最晚的待确认 ID以及各消费者分别积压了多少。如果某个消费者名下积压数量一直涨说明它处理不过来或者已经挂了。# 查看 pending 消息的具体 ID XPENDING order-events group_a - 10 # 查看 Stream 的整体信息 XINFO STREAM order-events XINFO GROUPS order-eventsXINFO GROUPS能列出这个 Stream 下所有消费组以及每个组的游标、pending 数量。排查消费组之间互相影响的时候特别有用。4. 消息可靠性ACK、PEL、CLAIM 这一套是怎么兜底的4.1 消息什么时候算真正“安全”很多初学 Stream 的人有一个误解以为消息写入 Redis 就算安全了。实际上在消费者组模型下“写入成功”只说明消息在 Stream 日志里还没有任何消费者认领它。消息真正“被安全处理”了要等到消费者回调处理完业务、执行 XACK、消息从 PEL 里消失。这套机制其实借鉴的是 Kafka 的 offset 提交思路但比 Kafka 更细粒度。Kafka 里消费者更新的是 partition 级别的位置而 Stream 的 PEL 是消息级别的哪些消息处理完、哪些没处理完Redis 一清二楚。代价就是 PEL 需要维护额外的内存消息量大且迟迟不 XACKPEL 会膨胀得很厉害。所以在线上用 Stream你一定要设一个监控XPENDING的数量如果持续增长就要告警。它相当于消费端的“积压水位”。4.2 消费者宕机了消息怎么捞回来假如消费者 A 读走了消息 MPEL 里记录下了 M但 A 在执行业务逻辑时宕机了M 就卡在 A 的 PEL 里。这时候 group 里其他消费者 B 是看不到 M 的因为 M 已经被投递给 A 了。怎么把 M 从 A 手里转移给 B用XCLAIM。# 将 group_a 中 consumer-1 的 pending 消息转给 consumer-2 XCLAIM order-events group_a consumer-2 60000 1716888888888-0第三个参数60000是最小空闲时间单位毫秒。意思是只有当这条消息在 A 的 PEL 里至少待了 60 秒没确认才允许转移。这个时间窗口是必要的因为你不确定 A 是不是还在慢慢处理万一 A 只是处理得慢你把消息转给别人就会重复处理。XCLAIM执行成功消息会从 A 的 PEL 移到 consumer-2 的 PEL同时消息的投递计数会增加。这里有个经验XCLAIM之后consumer-2 要用XRANGE或XREADGROUP指定消息 ID 把消息内容取出来重新处理不能只是 claim 完不管。后来 Redis 6.2 加入了XAUTOCLAIM能自动扫描一个消费者名下超时的 pending 消息并批量转移比手动 XCLAIM 省事很多。如果你用的版本高优先用XAUTOCLAIM。4.3 重复消费不可避免幂等设计才是兜底一个必须想清楚的问题Stream 能保证 at least once不能保证 exactly once。消费者 crash 之后被 claim 重放消息一定会被处理两次甚至更多次。很多人把重点放在“怎么让 Redis 不重复投递”上方向就错了。Redis 能做的只是在消息投递状态上尽量精确但只要你的业务处理有网络超时、进程重启这类不确定因素重复就在所难免。所以真正的解法是消费端幂等要么用业务流水号去重要么把 Stream 消息 ID 存到数据库唯一索引里要么让下游操作天然幂等比如“把状态置为已支付”这种覆盖式写入。我在实际项目里的做法是每条消息带一个业务幂等键消费端拿到消息后先查 Redis 本身的去重集合处理成功再把幂等键写入集合并设置和业务超时相匹配的过期时间。这样即便消息被重复投递幂等键也能挡住。5. 持久化、内存管理与性能调优的实践心得5.1 持久化配置对 Stream 的影响别忽略Redis 本质是内存数据库Stream 的消息也住在内存里。那消息的安全性靠什么靠 RDB 快照和 AOF 日志。RDB 是定时做全量快照如果 Redis 在两次快照之间宕机快照之后写入的 Stream 消息会丢。AOF 是追加日志更细粒度地记录写操作。你可以在redis.conf里调整appendfsync的策略always每个写命令都刷盘最安全但性能下降明显。everysec每秒刷一次最多丢一秒数据性能折中生产用的最多。no交给操作系统刷盘可能丢更多数据。如果你拿 Stream 做核心业务消息队列我的建议是把appendfsync设为everysec同时打开 RDB 作为兜底快照。还要设置合理的maxmemory策略防止消息堆积把 Redis 内存打爆。有人会说“既然 everysec 最多丢一秒数据那是不是还能丢消息”确实是。所以 Stream 更适合容忍少量丢失、但需要快速处理和高吞吐的业务比如活动削峰、异步通知、日志管道。真要强一致还是用专业 MQ 吧。另外一个容易踩的坑AOF 重写期间 Redis 会根据当前 Stream 状态生成紧凑的追加日志这个过程如果 Stream 特别大会占用额外内存和 CPU。建议在流量低谷期做重写或者配置自动重写阈值时留足余量。5.2 Stream Key 的内存膨胀和控制手段因为 Stream 是全量存内存无限追加意味着无限膨胀。所以生产环境一定要配合裁截策略。# 保留最近 1000 条消息超过的删掉 XTRIM order-events MAXLEN ~ 1000~是近似裁截意思是 Redis 不需要精确到 1000 条可以留多一点换取更高的效率。如果去掉~Redis 每次写入都会精确清理虽然内存控制精确但写入吞吐会受影响。我一般在写多读少的场景用MAXLEN ~配合一个略低于报警阈值的长度值。还要注意XDEL能删除指定消息但删除大段历史消息后底层 Radix Tree 节点可能不会立刻完全释放内存。如果 Stream 里堆积过大量消息、删除后内存一直没有回到预期水位可以等低峰期把 Stream 数据迁移到一个新 key 然后删除旧 key。5.3 集群模式下使用 Stream 的局限性Redis Cluster 对 Stream 的支持是基于 key 粒度的一个 Stream key 落在某个 slot 上由某个主节点承载。你做不到像 Kafka 那样把一个 topic 的分区均匀分布到多个 broker 上。这意味着当单个 Stream key 的写入或消费吞吐达到单节点瓶颈时扩展方案只有一个业务侧把数据拆分到多个 Stream key比如按用户 ID 尾号、订单号 hash、业务地域维度分片然后让消费者分组订阅多个 key。这种方法能横向扩展但也带来一个问题跨 key 的顺序性无法保证。假设同一个用户的消息落到了 key1 和 key2那这两条消息的执行顺序就不确定了。所以分片设计时要把需要严格有序的维度放到同一个 key 内。比如订单事件按订单号 hash同一个订单的所有事件都进同一个 Stream key顺序就保住了。这个限制在面试里一定要能讲出来因为很多人以为 Stream 和 Kafka 一样集群天然并发实际上完全不是一回事。6. 常见面试追问与线上实战避坑6.1 高频追问 Top 5我建议你怎么答问Stream 和 Pub/Sub 的本质区别是什么答Pub/Sub 是广播模式消息不落地消费者不在线消息就丢了Stream 是持久化的日志结构消息在 Redis 里保留消费者可以按 ID 回溯读取而且支持消费者组。一句话总结Pub/Sub 是瞬时的扇出Stream 是可靠的多播。问Stream 会不会丢消息答取决于你的持久化配置和消费方式。如果只开 RDB 不开 AOF宕机会丢最近的数据开了 AOF everysec极端情况丢一秒。使用消费者组并且正确 XACKRedis 能在消息投递状态上尽量精确但消费端崩溃后通过 XCLAIM 恢复时可能重复投递。所以消费逻辑要做幂等。问消费者组里加一个新的消费者能消费到历史消息吗答不能自动消费到历史消息。组的游标是独立维护的新消费者加入后只能消费组游标之后的新消息。如果想让新消费者从头消费历史消息可以另外创建一个消费者组指定起始 ID 为 0。记住同一份 Stream 数据可以创建多个消费者组组与组之间的游标完全独立。这有点类似 Kafka 不同的 group 消费同一份数据互不影响。问如果 Stream 消息堆积了很久怎么知道哪条没消费答用XPENDING查组内 pending 数量和具体 ID 范围用XRANGE按 ID 区间读消息内容。如果 pending 里大量消息超时用XCLAIM转移给其他消费者或者临时加消费者分摊。问什么时候不应该用 Stream答需要 exactly once、需要海量消息堆积、需要跨地域多活强一致、需要复杂消息路由这些场景 Stream 都不合适。它是轻量级可靠消息队列不是分布式消息中间件。6.2 线上用过之后我才知道的几个细节第一个坑是创建消费者组时的起始位置。很多人默认用0结果一创建就消费全量历史消息造成大量无效处理。如果只关心新消息要用$但$在创建那一刻就定下来了创建之后生产的新消息才会进入该组。还有XGROUP CREATE支持MKSTREAM如果 Stream 不存在可以自动创建空流免得还要先 XADD 一条消息。第二个坑是 BLOCK 参数别设 0。XREADGROUP的BLOCK 0表示永久阻塞如果消费者代码没做超时控制一旦 Redis 连接异常线程就会卡死在等待上。我习惯设置成 1000 到 3000 毫秒的超时循环里继续读取。处理完一批以后主动退出阻塞这样还能顺带做线程中断检查和心跳上报。第三个坑是关于消息内容的设计。Stream 的 field-value 都是字符串复杂结构要提前序列化。我见过有人直接塞 JSON 字符串消费端解析起来不方便还浪费空间。建议把通用字段比如事件类型、业务幂等键放在 Stream 的 field 里便于在 Redis 侧就用XINFO快速查看把业务数据整体作为 JSON 放在一个 field 里。这样调试的时候不用把整条消息拿出来解析一遍。6.3 我这一套在实际项目里的落地方案简单分享一个我最近做的订单异步通知场景。订单服务在创建订单时把order-created事件 XADD 到order-events这个 Stream key 里filed 包含order_id、user_id、timestampvalue 是序列化后的订单快照。下游有三组消费者一组负责发送短信通知一组负责积分变更一组负责给数据分析管道喂数。各组之间游标独立互不影响每组可以独立回溯消费。每组消费者部署了两个实例用同一个 group 名、不同的 consumer 名。实例启动时注册一个定时任务每 5 秒调用XPENDING检查自己消费者名下有没有超时未确认的消息有的话就重新拉取处理。处理成功的消息统一XACK处理失败超过三次的进入死信逻辑把消息 ID 和失败原因写到另一个 Redis 列表里人工补偿时再读出来。这个方案支撑了日均千万级的事件量Redis 内存峰值控制在 3GB 以内。如果你也是从 0 到 1 搭轻量消息队列这套结构可以直接抄作业。最后说一点真心话Redis Stream 是我觉得 Redis 生态里被低估的能力之一。它没有 Kafka 那样庞大的生态也没有 RabbitMQ 那么多路由协议但它足够简单、足够轻、足够快。很多中小团队根本没有必要为了一个削峰场景上全套 MQStream 用一个 Redis 实例就能扛住。关键是你要理解它的边界在哪里别拿它当万能药。把底层的 Radix Tree、消息 ID、PEL 这套机制吃透了你在任何场景里做取舍都会清晰很多。
返回列表