ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Redis Stream核心机制与工程实践全解析

Redis Stream核心机制与工程实践全解析 前阵子接手一个老项目的线上维护半夜告警群突然刷屏Redis 内存持续走高业务方反馈通知消息有延迟甚至有消息丢失。翻代码一看队列用的是 List BRPOP 那套经典方案生产端 LPUSH消费端 BRPOP看起来没什么问题但一压流量就暴露了本质缺陷——消息没有确认机制消费端一崩就是永久丢失也没有消费组的概念多实例扩容时根本没法把消息均匀分给每个 worker。那个周末我蹲在屏幕前把 Redis 文档翻了个底朝天最后把方案整体迁到了 Redis Stream 上。这是 Redis 5.0 引入的数据类型很多人到现在还停留在“听说过”的阶段实际上它解决的就是上面这类有状态、需要确认、需要消费者组协作的消息场景。这篇文章我想把 Redis Stream 从原理到实践讲透聊聊它到底解决了什么问题、适合什么场景以及落地时会踩到哪些坑。1. 为什么需要 Redis Stream—— 从一条消息说起1.1 从 List 队列到 Pub/Sub消息系统演进的两次转折在 Stream 出现之前Redis 做消息队列主要有两套姿势各有各的疼。第一套是 LPUSH BRPOP 这种列表队列。List 本身是双向链表左进右出天然就是一个 FIFO 队列。这套方案实现极其简单三四行代码就能跑起来直到今天仍然有很多老项目在用。但它的硬伤非常明显消息弹出之后就从结构里删掉了消费者拿到消息之后如果还没处理完就崩溃这条消息就再也没有人知道了。没有 ACK、没有重试、没有死信一旦业务代码里忘记处理某个异常分支数据就是静默丢失。另一个问题是消费端是“你争我抢”的模式多个 worker 同时 BRPOP 同一条队列谁抢到算谁的没法做到消息分片或者按组隔离。第二套是 Pub/Sub也就是发布订阅。它解决了解耦的问题生产者把消息扔到 channel 上所有订阅者都能收到一份。但 Pub/Sub 是“即焚”模式消息发出去之后如果没有消费者在线这条消息就消失在空气里Redis 不会为任何订阅者缓存消息。这带来的直接后果是只要消费者重启一下、网络抖动几秒钟你就永远错过了那几秒内产生的消息。再加上 Pub/Sub 不支持消息确认、不支持回溯它本质上更适合做实时通知、事件广播这类“丢了也不心疼”的场景而不是做业务消息队列。1.2 Stream 到底解决了什么问题三类需求的交叉点把上面两套方案的痛点放在一起看你会发现在很多业务场景里我们需要的是这样一个东西消息能持久化存下来生产端发完不用关心消费端是否在线消费端能确认消息已经处理完成没确认的消息可以被重新投递多个消费者可以组织成“组”协同消费同一条消息只被组内一个成员处理消费进度可以记录新加入的消费者能从最早或最新的位置开始读单个消息有稳定 ID方便查重、回溯、做补偿。这些需求单拎出来其实都对应着成熟消息中间件的能力。但很多团队没有到必须引入一套独立 MQ 的程度——运维成本、部署成本、团队学习成本都是实实在在的负担。Redis Stream 的价值恰恰出现在这个交叉点上它用 Redis 一个数据类型的能力把上面这些能力全部装了进去而且 API 设计得足够直观学习曲线比 Kafka 平缓得多。用一句话概括Redis Stream 让 Redis 从“缓存工具”变成了一种轻量级消息中间件。它并不试图取代 Kafka 或 RabbitMQ但在中低吞吐、偏业务内聚、不想引入额外组件的场景里它是一个非常合理的选择。2. 核心机制与关键概念2.1 数据结构本质一个可持久化的追加日志Stream 的本质是 append-only log中文叫追加式日志。所有消息按照写入顺序依次追加到 Stream 尾部每条消息有一个全局唯一 ID结构上类似一个只允许追加的列表但内部实现比 List 复杂得多使用基数树来索引消息 ID所以历史回溯的效率高得多。操作层面它和 List 最大的区别是List 的 BRPOP 会把元素从结构里移走而 Stream 的 XREAD 只是读取消息仍然留在 Stream 里。这样消费速度慢不会导致消息丢失你随时可以从任何位置重新读。这也是它和“队列”最本质的区别——Stream 本质上是日志日志天然允许你反复读取、回溯、补数据。Redis 官方把它称作 Stream 而不是 Queue是有意的命名选择。队列的含义是“管道的另一端有人等着拿”而日志的含义是“事情发生了我就记录下来”。谁读、读到哪里、读完怎么处理这些完全是消费者自己的事情。2.2 ID 机制时间戳 序列号的精妙设计每条消息的 ID 格式是millisecondsTime-sequenceNumber比如1711108800000-0。前半段是 Redis 服务器本地时间戳毫秒后半段是同一毫秒内的自增序号从 0 开始。这种设计有两个很漂亮的特性第一个特性是 ID 天然有序。因为 ID 是单调递增的用 XRANGE 按 ID 范围扫数据就是按时间顺序扫不需要额外的排序字段排查问题的时候直接按时间窗口拉一段数据出来看非常方便。第二个特性是客户端可以指定 ID。XADD 的时候你可以手动传入 ID只要比 Stream 里当前最大 ID 大就能插入。这在做数据迁移、重建数据的时候特别有用——你可以把旧系统的消息按照原来的 ID 顺序导入新 Redis保持全局消息 ID 的一致性。如果你希望 ID 完全由业务生成也可以用XADD mystream 1711108800000-0 field value这种形式指定但除非有强需求一般不建议这么做因为时间戳部分不是真实插入时间的话后续按时间回溯就会不准。2.3 消费者组协作消费、ACK 与 PEL 的三角关系Stream 最核心的能力是消费者组Consumer Group。消费者组是这样一个模型一组消费者共同消费一个 Stream组内每条消息只会被一个消费者领取领取之后进入该消费者的 PELPending Entries List待确认消息列表消费者处理完之后发送 XACK消息才会从 PEL 里移除。这里最关键的概念是 PEL。之前说 List 方案消息会丢就是因为没有这个结构。只要有消费者从 Stream 里读走了消息但还没 XACK这条消息就会一直躺在 PEL 里。就算消费者崩了、网络断了、进程重启了消息也还在 Redis 里记着账。你随时可以用 XPENDING 查看哪些消息没被确认用 XCLAIM 把超时未确认的消息转移给另一个消费者重新处理。消费者组的另一个优势是消费进度的持久化。每个消费者组在 Stream 上都维护自己的游标记录着“这个组已经消费到了哪条 ID”。因此同一个 Stream 可以挂多个消费者组各组之间互不干扰进度互不覆盖。比如一个订单系统同一份订单事件可以同时被“订单状态同步组”和“数据分析组”消费各自维护各自的消费进度这在实际业务里非常常用。3. 实操从零搭建一个 Stream 消息队列3.1 基础命令与一个最小闭环先看最基础的操作。假设我们要做一个工单创建通知生产端写入消息# 写入一条消息字段可以自由定义类似一个小 Hash XADD ticket:events * action create ticket_no T1001 user_id 9527 1711108800000-0*表示让 Redis 自动生成 ID返回的 ID 就是这条消息的全体ID。如果写入成功说明这条消息已经持久化到 Stream 里了。读取消息用 XREAD# 从 Stream 头部开始读 XREAD COUNT 10 STREAMS ticket:events 0这里0表示从最小 ID 开始读也就是从头消费。如果业务上只需要读最新消息可以把0改成$表示只读取调用时刻之后新写入的消息。最小闭环还要包括长度查看和范围查询 XLEN ticket:events XRANGE ticket:events 1711108800000-0 COUNT 10XRANGE 的参数是起始 ID 和结束 ID-和分别代表最小和最大。这套命令组合起来你可以随时按 ID 段把任意时间段内的消息拉出来排查这是 List 和 Pub/Sub 完全做不到的。3.2 消费组 ACK 的完整实现下面是实际项目里最常用的一套流程。先把消费组建出来# 创建消费者组从头部0开始消费 XGROUP CREATE ticket:events group_ticket 0如果想让组从创建时刻开始只接收新消息第三个参数用$如果要从头消费存量消息用0。生产环境建议单独加一个参数MKSTREAM组创建时如果 Stream 还不存在会自动建空 Stream避免因为“还没有任何消息”导致创建失败XGROUP CREATE ticket:events group_ticket $ MKSTREAM消费者读取消息用 XREADGROUP XREADGROUP GROUP group_ticket worker1 COUNT 10 BLOCK 5000 STREAMS ticket:events 这里几个参数值得逐个解释GROUP group_ticket worker1指定用哪个组、以及当前消费者的名字。组内每个消费者的名字必须唯一Redis 用消费者名字来登记 PEL名字一乱消费记录就乱。COUNT 10一次最多读取 10 条。这是批量拉取的意思可以减少网络往返。BLOCK 5000如果暂时没有消息阻塞等待最多 5 秒超时后返回空结果。不设 BLOCK 就变成非阻塞读有消息立即返回没有消息立即返回空。这个特殊符号是关键它表示“从组当前游标之后读取新消息”。如果不写你也可以写一个具体的 ID表示从这条 ID 之后读消息但语义就变成了“查看 PEL 里的历史消息”和正常消费新消息是两回事。拿到消息后业务处理完毕必须发送确认 XACK ticket:events group_ticket 1711108800000-0 1711108800000-1XACK 是把消息从当前消费者 PEL 里移除的唯一方式。如果忘了 ACKRedis 会一直认为这条消息“正在被处理”XPENDING 里的数字只增不减。用 Python 写一个消费循环组合起来大概是这样import redis import json r redis.Redis(hostlocalhost, port6379, decode_responsesTrue) stream ticket:events group group_ticket consumer worker-1 while True: # 阻塞读取最多等 5 秒 resp r.xreadgroup( group, consumer, streams{stream: }, count10, block5000 ) if not resp: continue for stream_name, entries in resp: for msg_id, fields in entries: try: # 业务处理比如推送通知、更新工单状态 handle_ticket(fields) # 处理成功确认消息 r.xack(stream, group, msg_id) except Exception as e: # 失败则不做 ACK让消息留在 PEL 中等待重试 logger.error(fhandle message failed: {e}, msg_id: {msg_id})这段代码是 Stream 消费端的基本范式读取、处理、确认三件套。注意异常分支里没写 XACK这是刻意为之——消息处理失败时留在 PEL 里后面才有“重新投递”的依据。3.3 横向扩展消费者增加时消息如何分配消费者组最实用的场景是横向扩容。假设 group_ticket 组最初只有 worker1 一个消费者它消费的进度游标可能在 1000。这时你新起一个 worker2 加入同一个组Redis 不会自动把游标分一半给 worker2而是采用一种“有状态分配”的策略。具体来说当一个新消费者加入并且用发起了第一次读取请求时Redis 会尝试把组里那些“已被分配但还没 ACK”的消息也就是 PEL 里的条目转移一部分给它。如果 PEL 里没有任何待处理消息新消费者会从组当前游标处继续读新消息不会回头补旧消息。这在扩容时意味着如果存量消息都已经消费完毕并 ACK新增的 worker 只会承担“未来新消息”的一部分如果存量积压非常严重新消费者加入后会被自动分配一部分 PEL 里的积压消息分担组内压力。不过坦白说Redis 待处理消息的自动重新分配策略毕竟不是 Kafka 那种精确的分区均衡它没法保证“每个消费者手里的消息数量完全一致”只保证“同一条消息不被同组两个消费者同时拿到”。在实际场景里如果积压量巨大我一般建议先把组内消费者的消息都消费完、清空 PEL再在低峰期扩容或者直接创建一个新的消费者组从存量位置重新消费用临时代码做数据补齐避免让 Redis 的分派机制处理极端积压。3.4 内存控制XTRIM 与消息过期策略Stream 的持久化特性带来的副作用就是内存增长失控。Stream 本身没有 TTL 概念所有消息会一直常驻内存如果生产端写入频率高又没有上限约束Redis 内存迟早被打爆。这就是我开头说的那个老项目告警的常见根因。控制长度的命令是 XTRIM# 只保留最近 1000 条消息 XTRIM ticket:events MAXLEN 1000 # 按内存近似截断 XTRIM ticket:events MAXLEN ~ 1000MAXLEN后面可以直接跟精确数字但每次插入都精确修剪效率低。加上~之后Redis 会在合适的时机才做裁剪允许结果略超目标长度性能好很多。在 XADD 的时候也可以直接带上MAXLEN ~ 1000选项让 Stream 在写入时自动维持上限。XTRIM 还存在一个潜在的坑它只会删除消息不会自动处理 PEL 里那些还未确认的引用。如果某条未 ACK 的消息被 XTRIM 删掉消费者试图 XCLAIM 或重新处理时会发现消息已经不存在。所以设置 MAXLEN 时要考虑好这个问题要么消费速度足够快PEL 里停留时间很短要么消费组读取频率高到旧消息不会积压到被裁剪。否则你需要在业务层面对“读到已被删除的消息”做容错别让程序一发现消息不存在就直接崩溃。4. 可能被忽视的细节阻塞、持久化与可靠性4.1 阻塞读的正确理解BLOCK 并不是死等XREAD 和 XREADGROUP 的 BLOCK 参数很容易被误解成“一直阻塞到有消息为止”。实际上BLOCK 单位是毫秒它的完整语义是“最多等这么久”超时后立即返回空结果。比如BLOCK 5000表示最多等 5 秒。如果没有消息5 秒后返回(nil)或者空列表客户端需要自己决定是继续循环还是退出。为什么要有这个超时设计因为客户端与 Redis 之间的长连接可能在空闲时被网络设备断开。如果客户端无限期阻塞读一旦连接被中间设备回收客户端实际已经收不到消息了但应用进程不知道会一直挂在那里。我在实际项目里见过不止一次这种情况某个 worker 进程看起来还在但已经“假死”阻塞读的 socket 早就断了Redis 端也检测不到直到手动重启才发现消息积压了十几万条。所以我的经验是生产环境里的 BLOCK 时长不要超过 10 秒配合客户端的重连逻辑来做。比如综合上面的 Python 例子把 BLOCK 设为 3000~5000 毫秒就是一个比较稳妥的选择。即使连接断了最多 5 秒后调用就会返回你可以在循环里捕获连接异常、重新建立连接、继续消费整个过程对业务无感。4.2 持久化真相Stream 数据到底安全吗Redis 的持久化方式决定 Stream 的安全级别。如果只开了 RDB 快照持久化Redis 进程崩溃后可能会丢失最近一次快照之后写入的数据Stream 也不例外。如果开了 AOF并且appendfsync配置是everysec极端情况下也会丢最多 1 秒的数据。只有appendfsync always才能在理论上做到每条命令都落盘但代价是写入性能骤降。很多团队把 Redis 当成纯缓存持久化配置很随意。一旦把 Stream 当消息队列用就必须重新审视持久化配置。毕竟消息队列丢消息和缓存丢数据是性质完全不同的事。我踩过的一次教训是生产环境 Redis 配置里 AOF 没开某天 Redis 实例 OOM 触发重启Stream 里积压的所有未消费消息瞬间全部蒸发。业务方反馈“消息层数据对不上账”最后只能靠数据库日志人工补偿。所以给一个明确建议如果 St民ream 里有不可丢失的业务数据务必开启 AOF并且把appendfsync配置为everysec甚至always同时给 Redis 配好内存上限maxmemory以及淘汰策略比如noeviction避免 OOM 直接杀掉进程。4.3 主从与哨兵模式下 Stream 的边界情况使用主从复制或哨兵模式时Stream 的消费行为和普通数据复制其实是一致的消费者只能从主节点写、从主节点读从节点只做数据备份和故障切换。这里有几个容易出问题的点。第一个是 Redis 故障切换后消费者的 PEL 记录是否存在。答案是存在的——因为消费者组的元数据也存在 Redis 里并且会复制到从节点切换之后组信息和 PEL 都还在。但前提是切换前这些元数据已经被复制到从节点。如果恰好主节点崩溃前还没来得及把最新的 PEL 状态同步到从节点那些更新就会丢消息会被重复消费。第二个是客户端连接的拓扑关系。Redis 客户端通常需要配置哨兵或集群模式下的“只在主节点读写”策略。如果消费者误连到了只读从节点XREADGROUP 会直接报错因为消费者组的创建和读取都要求在主节点执行。4.4 Stream 与“连接中断”网络错误到底是谁的锅很多人把 Stream 使用中的网络报错误认为是 Stream 本身的问题常见的有两类信息一类是stream disconnected before completion: transport error: network error另一类是stream disconnected before completion: idle timeout waiting for sse。这里的stream实际上指的是网络数据流TCP/HTTP 流式数据传输跟 Redis 的 Stream 数据类型没有任何关系。但在使用 Redis Stream 时类似的现象确实可能出现客户端做阻塞读长时间空闲之后连接被防火墙、云负载均衡器或 Redis 服务端回收客户端读操作超时或报连接重置。这种问题的本质是长连接空闲超时而不是消息队列故障。处理办法其实很简单客户端开启 TCP keepalive 参数XREAD 的 BLOCK 时间缩短消费循环里捕获连接异常后自动重连在 Redis 服务端调大timeout配置或者在客户端连接池配test_on_borrow之类的连接健康检查。把这些兜底做扎实之后网络层面的抖动就不至于导致消费线程挂死。5. 常见问题与排查实战5.1 线上消费停滞为什么 PEL 一直有未确认消息消费停滞是 Stream 使用中最高频的问题。症状是消息不断进来但业务处理延迟越来越大用XPENDING一看某个消费者的 PEL 里堆积了大量未 ACK 的消息。排查思路是先确认消费者进程是否还活着。我遇到过的情况是消费者进程活着但业务依赖的下游接口变慢每条消息耗时几十秒自然消费速度跟不上。这时候优先解决下游问题消息本身没有丢等下游恢复之后 PEL 会被慢慢清空。另一种情况是消费者进程假死就像前面说的阻塞读连接断开。这种必须靠消费循环里的超时 重连逻辑来兜底否则你可能要手动杀掉进程才能恢复。排查工具我推荐看这几个# 查看组里每个消费者的积压情况 XPENDING ticket:events group_ticket # 查看某个消费者 PEL 里最早未确认的消息 ID XPENDING ticket:events group_ticket - 10 worker1 # 查看 Stream 最新消息 ID 和当前组游标位置判断组是否落后 XINFO GROUPS ticket:events XINFO STREAM ticket:eventsXINFO 输出里lag字段表示组的消费进度落后了多少条。如果 lag 持续增长说明消费端处理能力跟不上写入速度优先扩容消费者其次检查业务逻辑。5.2 消息重复消费从原理到幂等设计Stream 的消息重复不是 bug而是设计如此。消费者在 ACK 之前崩溃PEL 里那条消息会一直被标记为未确认。重启后如果用 XAUTOCLAIM 或 XCLAIM 重新处理这条消息就会再次进入消费者手里。所以 Stream 交付语义是 at least once至少一次不是 exactly once精确一次。这意味着使用 Stream 的业务代码必须设计成幂等的。幂等的含义是同一消息处理两次和执行一次最终效果完全一致。举个例子处理工单创建事件时不要无脑 INSERT而是先根据消息 ID 查一下工单是否已存在或者数据库表字段加唯一索引处理余额变更时不要把“增加 100 元”做成一个无条件 UPDATE而是把事件 ID 作为去重键用一条带条件的 SQL 或 Redis SETNX 保证同一事件只生效一次。我在团队内部定的规矩是写 Stream 消费者时第一件事就是设计消息去重而不是先实现业务处理。去重键一般用消息 IDStream ID或者业务里的唯一流水号两者都行只要实现简单可靠。5.3 死信与补偿处理永远失败的消息任何消息队列都有消息重复消费很多次仍然失败的情况。Stream 本身没有内置的死信队列你需要自己设计一个变通方案。我的做法是在消息字段里附加一个自定义的重试计数字段比如retry_count。消费者读取消息后如果处理失败调用 XADD 往一个专门的ticket:events:deadStream 里写一条同样的消息带上原始消息 ID 和当前重试次数同时 XACK 掉原消息避免它一直在主 Stream 的 PEL 里挂账。之后由另一个定时任务扫描死信 Stream按重试次数决定是直接人工处理、还是退避后重新投递。核心原则是不要让同一条消息无限重试也不要让失败消息永久占着 PEL 里的内存。设计好重试上限和死信迁移流程线上才不至于被阻塞消息拖死。5.4 客户端连接异常与 Redis 端超时参数的平衡前面提过网络抖动会导致消费者连接断开。如果 Redis 配置了timeout 300300 秒空闲断开那么一个消费者阻塞读超过 300 秒时连接就会被服务端主动关闭。阻塞读的 BLOCK 参数如果设得比服务端 timeout 还大就很容易出现“读得好好的突然连接断”的假象。解决原则也简单客户端 BLOCK 必须小于服务端 timeout。一般我习惯把服务端 timeout 设为 0也就是永不断开空闲连接靠客户端自己管理连接生命周期如果公司安全策略不允许至少确保 BLOCK 时长 timeout / 2留出足够的重连窗口。6. 选型与边界什么场景用它什么场景别用6.1 与 Kafka、RabbitMQ 的横向对照很多团队选型时会纠结Redis Stream 和 Kafka、RabbitMQ 到底怎么选。我画一张简表对比一下实际使用感受维度Redis StreamKafkaRabbitMQ部署成本低Redis 本身就在高要搭 ZK/Broker 集群中独立 Erlang 节点吞吐量上限中受单机内存和 CPU 限制非常高分区并行写入中复杂路由场景下表现好消息确认XACK 机制较灵活Offset 提交机制手动/自动 ACK消费组有但分派策略较朴素成熟的分区再均衡成熟的竞争消费模型消息回溯XRANGE 按 ID 范围灵活按 offset 或时间支持好一般死信/延迟队列无内置需自建有延迟队列思路需配置内置死信和延迟队列适用规模单机/主从几十万 QPS 内海量日志/事件流复杂路由、业务系统这张表不是用来分高下的。实际判断标准应该看你的业务约束团队有没有运维 Kafka 集群的人力消息量真的需要分区扩展吗需要 Topic 级别的海量留存吗如果答案都是“不太需要”那 Redis Stream 完全够用而且省心很多。6.2 建议用它和别用它的场景我个人经验里Redis Stream 非常适合这几类场景业务内部的异步解耦比如订单创建后发消息给通知服务、积分服务、审计服务各自挂一个消费者组小型数据同步管道比如从一个 MySQL 表同步变更到另一个服务数据量不大但要求不能丢需要“临时低下游”的系统消息先落在 Stream 里下游恢复后再消费天然带缓冲快速原型和中小团队项目不想引入额外中间件Redis 已有能力够用。反过来这几类场景建议直接上 Kafka 或 RabbitMQ日写入量上亿级别或者需要跨机房、多副本高可用保障需要保留消息多天甚至数周、按时间维度做大规模离线回放分析需要复杂的消息路由、延迟队列、死信队列、发布订阅等企业级特性团队规模大有专职中间件运维人力能把 Kafka 的存储和分区管好。Redis Stream 的单机写性能再强也有物理上限水平扩展也不是它的强项。它最大的卖点是“轻”和“够用”而不是“强”和“全面”。最后再分享一个在实际项目里的小技巧。我在消费端代码里总会给 PEL 加一个定时巡检逻辑用 XPENDING 每 30 秒检查一次各组积压情况。如果某个消费者 PEL 里的消息数超过阈值并且消息的最早时间戳已经超过 5 分钟就说明该消费者大概率出了问题这时候报警提示人工介入而不是依赖 Redis 自身的机制去“自动恢复”。消息队列这东西稳定性终究要靠外围的监控和纪律来保障。Stream 给了你足够好的底子但真正让它跑稳的还是使用它的人对细节是否足够较真。
返回列表