凌晨三点,我被一条“消息积压”的监控告警从梦里拽起来。打开面板一看,积压的倒不是队列本身,而是一批已经被消费者读走的订单任务,在消费者进程崩溃之后,像从来没有存在过一样。那是一个用得很“经典”的List队列方案:生产方LPUSH,消费方BRPOP,简单、高效、人人都这么写。可也是那一次事故之后,我才把Redis Stream这个从Redis 5.0引入后被我冷落很久的数据结构翻出来,从源码到命令到落地场景完整磕了一遍。这篇文章不打算只讲API怎么用,而是想认真聊一个更底层的问题:Redis Stream到底解决了什么问题,以及我们在架构选型时,该怎么判断“什么时候该用它,什么时候该去上Kafka”。
如果你正在选型团队的异步任务队列、事件总线或审计日志存储,如果你觉得Redis的List队列“用着还行但总有点慌”,这篇文章应该能给你一个比较完整的答案。
1. 一个让我重新审视Redis消息体系的凌晨告警
1.1 事故现场:为什么“队列积压”背后其实是“消息失踪”
那套系统的业务很简单:用户下单后,订单服务把一条任务塞进Redis队列,一个Java后台线程池从队列里取任务,然后调用库存、积分、通知三个外部服务。大概跑了半年,一切正常。直到某个晚上,业务方反馈说有几笔订单支付回调没触发积分发放,排查下来发现:消费线程所在的容器在凌晨被健康检查强制重启了,重启前它已经从队列里BRPOP出来一批消息,但还没执行完业务逻辑,进程就没了。Redis里的List是一个“取出即删除”的结构,BRPOP返回的那一刻,数据就从Redis里消失了。进程还没来得及处理,消息就永久丢失。
监控面板上只显示队列长度,而这批消息已经被取出,所以队列长度是0。我看到的“积压告警”其实是另一个索引队列的消费速度下降触发的,真正的损失消息却完全不可见。那一刻我意识到,一个号称“消息队列”的方案,居然连“投递成功的消息是否被处理成功”都回答不了。
1.2 这次事故暴露的三个核心问题
这次事故其实暴露了三个几乎所有业务团队早晚会撞上的问题。
第一,消息被取走但没被处理,系统怎么识别“消费失败”?List结构根本没有“投递状态”这个概念,BRPOP返回就默认“你处理了”,后面发生什么都不归Redis管。第二,多个下游服务需要各自消费同一条事件时,怎么独立推进?比如订单创建这件事,风控要看、积分要加、短信要发,它们消费速度不同,同一个List里的一条消息被一个消费者抢走后,其余服务就看不到这条消息了。第三,消息处理完之后,能不能按时间回溯?“根据订单号查一下这条消息当时长什么样”这类审计需求,List基本做不到,等你想到要查的时候,数据早就pop走了。
没有哪个问题是List“完全不能用”,但三个问题叠在一起,Redis Stream这个方案就变得非常必要了。它把“可靠投递”“消费组”“可按ID回溯的日志”这三件事,在Redis内部做成了原生的数据结构能力,而不需要业务层再叠一堆临时表、备份列表、定时任务去修补。
2. Stream出现之前:三套Redis“思路”,以及它们各自的边界
2.1 List的LPUSH/BRPOP:只适合做一个“简单待办队列”
List队列可能是Redis世界里最常见、最朴素的消息用法:一个生产者不停LPUSH写左侧,多个消费者BRPOP从右侧取。这套模型的优点是吞吐高、延迟低、实现成本几乎为零,几个命令就能跑起来。它的缺点也很明确。
首先,在标准消费流程里,一条消息只能被一个消费者get到,也就是说它本质是“任务分发型”的,不是“广播型”的。同一个业务如果想要两个服务同时独立处理同一批消息,需要复制两份队列,自己控制分发逻辑。其次,没有确认机制,消费者取走消息后无论是否处理失败,消息都已出列。BRPOPLPUSH可以先把消息备份进一个备用List,超时后再扫描回收,这算一种“手工可靠性”方案,但超时时间怎么定、重复处理怎么去重、多个消费者同时扫描会不会抢消息,都是问题。最后,它完全没有“消费进度”的概念,无法回答“我们处理到哪一条了”。
如果只是“后端任务随手异步一下”,用List没问题。一旦涉及订单、支付、资金类业务,把List当正式消息队列用,就是把整个业务的可靠性押在“消费者进程绝不宕机”这个假设上。
2.2 Pub/Sub:广播很开心,但订阅者必须“恰好醒着”
我见过不少团队把Redis Pub/Sub当成事件总线用,发布者PUBLISH,订阅者SUBSCRIBE,看上去很优雅。但它有一个致命前提:发布者发出消息时,如果订阅者不在线,这条消息就直接丢弃了。它没有持久化、没有ACK、没有积压概念,消费端哪怕断线一秒钟,这一秒钟的订单事件就永远追不回来。
所以Pub/Sub适合的场景其实是“实时性要求高、丢几条无所谓”的通知,比如服务在线状态广播、本地缓存失效通知、WebSocket节点间转发。你想让它承担订单事件这种核心业务链路,等于默认业务方可以接受消息丢失。另外,Pub/Sub的积压能力是在每个订阅者客户端的内存里的,一旦消费跟不上,积累在客户端连接上的消息可能导致连接阻塞、内存暴涨,甚至被Redis强制断开。
2.3 用ZSet或MySQL手搓队列:能排序,但“搓”不出可靠性
在Stream发布之前,还有人用ZSet模拟有序队列:score存时间戳,消息内容存member,消费者用ZRANGEBYSCORE取最早的消息,再用ZREM删除。这套方案能实现“按时间排序”和“回溯”,但ZREM同样没有确认机制,取出来就是删掉,消费者挂了依然丢消息。而且取消息和删消息不是原子操作,并发控制稍一疏忽就会重复消费。
更麻烦的是,ZSet方案里每一个能力——消费组、ACK、死信、偏移量——都要自己用事务或者Lua脚本去实现,实现完还要考虑过期清理、性能退化。说白了,这是用业务代码去造一个Redis内部本不该由你关心的轮子。Stream出现之后,这种手搓方案应该被收进历史垃圾桶了。
这三套方案代表了三类需求:List满足“分任务”,Pub/Sub满足“广播”,ZSet满足“排序”。但它们没有一个能同时满足“可靠、可分组、可回溯”。Redis Stream正是把这三条线整合进了同一种数据结构。
3. Stream核心模型:与其说它是一个队列,不如说它是一个“可共享指针的单机日志”
3.1 XADD:消息不是塞进队尾,而是追加进日志
我第一次接触Stream时犯了一个思维错误:总拿它跟List队列比,觉得“XADD就是LPUSH的替代品”。后来读文档才意识到,Stream的更准确类比是日志(log),不是队列。XADD往Stream里追加一条记录时,它可以自动生成一个ID,形式是<毫秒时间戳>-<序号>。比如1684123456789-0。时间戳部分保证消息大体按时间排序,同一毫秒内的多条消息通过递增序号区分。
这个ID是整个Stream模型的地基。因为消息ID天然有序,你可以通过XRANGE命令按范围查询任意一段消息:
# 查所有消息 XRANGE order:events - + # 查某个时间窗口内的消息 XRANGE order:events 1684123000000-0 1684124000000-0List队列里,消息从一端进、从另一端出,中间状态几乎不可见。而Stream里的消息一旦写入,默认就一直存在(除非你裁剪),你随时可以把时间轴拨回去重看。对于审计、排查、数据对账,这个特性非常值钱。
3.2 消费者组偏移量机制:组内一份,组间多份
Stream的消费者组,简单理解就是“每个消费组都维护了一个独立的游标”。同一个Stream里可以创建多个消费组,每个组都觉得自己在消费一个独立的队列,组与组互不干扰。
比如订单事件流order:events,可以有两个组:group:risk和group:points。风控组从第一条往后读自己的,积分组也从头往后读自己的,互相不影响。组内的多个消费者则共享这一份游标,Redis大致按轮转方式把新消息分给组内的消费者,保证同一条消息默认只被组内一个消费者读到。这一点和Kafka的消费者组语义高度一致。
这里的“>”符号也很关键。XREADGROUP读取时如果传>,代表“只读取该组从未投递过的新消息”;如果传具体ID,就代表“从这条ID开始重读”。后者是实现“重新消费”的钥匙,也是和普通队列拉开差距的地方。
3.3 PEL与XACK:系统如何判断“消息没处理完”
Stream最打动我的机制是Pending Entries List,也就是每个消费者组的待确认消息列表。消费者用XREADGROUP读走一条消息后,消息不会像List那样消失,而是被记入这个消费组的PEL,并绑定到具体消费线程名下。等业务处理成功,消费者调XACK通知Redis“这条我搞定了”,Redis才把它从PEL里移除。
如果消费者进程崩溃,还没来得及XACK,消息会一直留在PEL里。另一个消费者可以通过XCLAIM把这批超时未确认的消息转移过来继续处理。这个机制,直接把我经历过的那次“凌晨事故”里的“消息失踪”问题解决了:没有ACK的消息,不会凭空消失,只会在PEL里躺着等你处理。从数据可靠性角度讲,Stream相当于把“投递”和“确认”拆成了两个独立动作,而业务逻辑不需要自己做任何额外持久化。
3.4 为什么“日志型”结构更适合重放与审计
Stream之所以被设计成追加式日志,还有个很实际的原因:分布式系统里“消息到底长什么样”是事后排查的关键。List把消息pop掉之后,Redis里就干干净净,想取证都没得取。Stream则保留完整的时间轴,配合XRANGE可以迅速回答“某段时间内到底有哪些订单事件”“这条消息的内容字段是什么”。
另外,Stream的消息体是field-value结构,类似一个小型Hash,天然适合传输结构化数据。你不需要在外面再包一层JSON字符串来解析,虽然很多团队实际也会把JSON塞进value里,但结构本身是支持的。
4. 和Kafka的差别:Stream不是“小Kafka”,而是“内嵌在Redis里的可靠消息通道”
4.1 模型相似,但定位完全不同
网上很多文章把Redis Stream称作“小Kafka”,我觉得这个说法容易误导。它们确实共享了很多概念:消息ID/offset、消费者组、重放、持久化日志。但如果你拿Kafka的期望值去要求Stream,很快就会遇到瓶颈。
Kafka是一个独立的消息中间件集群,核心能力是“把海量数据流以极高吞吐持久化到磁盘,并通过多分区水平扩展”。Redis Stream则是Redis里的一个数据结构,它本质上跑在单机内存(或集群中某个key所在的分片)上,能给你的是“可靠的、可回溯的、支持消费组的消息通道”,但它的横向扩展能力远不如Kafka,数据保存量也受内存限制。
4.2 分区能力:一个Stream对应一个分片
这是两者最关键的差异。Kafka可以按key、按业务把一个Topic拆成几十个partition,每个partition独立有序,消费者组内按分区分配任务,整体吞吐随分区线性增长。Redis Stream呢?单个Stream就是单个有序结构,没有原生的分区概念。你可以手动按业务维度拆成多个Stream key,比如order:events:province1、order:events:province2,但拆分逻辑要自己设计,消费者端也要自己路由。
对大多数内部业务而言,单Stream的吞吐通常不是问题。单机Redis Stream压到十几万QPS写入是有可能的,很多业务的实际消息量也就每秒几百到几千。只有当你需要支撑百万级QPS、跨集群、跨机房的数据管道时,Stream的“单分片”模式才会变成明显的天花板。
4.3 持久化语义不一样
Kafka的消息是直接落在磁盘上的,所以它可以承载几十GB甚至TB级别的日志流,而且消息保留策略独立于应用进程。Redis Stream的消息虽然也能通过RDB/AOF持久化,但它是内存中的数据结构,消息量过大会把内存吃穿。你可以用MAXLEN或XTRIM限制Stream长度,但这就意味着历史消息一定会被裁剪。Kafka靠多副本和磁盘解决容灾,Redis则靠原生的RDB/AOF机制,本质上两者对“长期可靠存储”的投入不在一个量级。
这两个产品的选择,不是用“谁替代谁”来思考的,而是先回答三个问题:消息量级是每秒几百还是每秒百万?数据要保留多久,是几小时还是几个月?团队是否已经有一套Kafka集群?如果你的团队都还没引入Kafka,只是想快速解决“花钱买不了消息丢失”的问题,Redis Stream往往比搭一套Kafka划算得多。
4.4 一张表理清选型
| 维度 | Redis Stream | Kafka |
|---|---|---|
| 存储 | 内存为主,可裁剪 | 磁盘追加日志 |
| 分片能力 | 单key单分片 | 多分区水平扩展 |
| 消费者组 | 支持 | 支持 |
| 单条ACK | XACK原生支持 | offset提交,无单条语义 |
| 消息回溯 | XRANGE方便 | 任意offset |
| 吞吐量级 | 单实例十万级 | 集群百万级 |
| 部署复杂度 | 复用已有Redis | 独立集群 |
| 典型场景 | 业务任务队列/事件分发/审计 | 大数据管道/跨团队事件中枢 |
5. 它到底能解决什么:三类让我“值回票价”的应用场景
5.1 可靠的异步任务队列:核心场景,没有之一
我最推荐团队拿来开刀的场景,是把原先List队列的异步任务迁移到Stream。比如订单支付成功后要发通知、积分、写日志,以前串行做太慢,异步搁List里又怕丢。用Stream可以这么写:
# 生产端,写入一条订单事件 XADD order:events * order_id 10001 status paid # 创建消费组,0表示从第一行开始读 XGROUP CREATE order:events group:notifier 0 # 消费者A读取新消息 XREADGROUP GROUP group:notifier worker-A COUNT 10 STREAMS order:events >消费端的核心逻辑是:读消息、处理业务、显式XACK。
import redis r = redis.Redis(host="...", decode_responses=True) while True: resp = r.xreadgroup( groupname="group:notifier", consumername="worker-A", streams={"order:events": ">"}, count=10, block=5000, ) for stream_name, messages in resp: for msg_id, fields in messages: try: handle_notify(fields["order_id"]) r.xack("order:events", "group:notifier", msg_id) except Exception: # 不XACK,消息留在PEL里,后续由XCLAIM补处理 log.error(...)这套流程带来的收益是实打实的:业务处理失败时,消息不会消失;消费者重启后,PEL里积压的未确认消息可以被重新认领;如果某个worker卡死了,其他worker通过XCLAIM把它的消息转移走再处理。相比List方案,可靠性直接上了一个台阶。
5.2 让多个下游以各自速度消费同一份事件
微服务架构里最烦的事之一,就是同一个业务事件要被多个团队消费。订单创建了,风控服务要算风险,积分服务要加积分,短信服务要发通知。它们各自的处理速度和失败率都不一样,如果把事件塞进同一个List,只能有一个服务消费;塞进Kafka,又得专门搭一套基础设施。
Stream的消费者组天然解决这个问题。每个下游服务建一个自己的消费组,各自从Stream里独立读事件,互不干扰,还可以各自维护自己的消费进度。某个服务挂了,不影响其他组继续消费。这个“组间独立”的特性,让Stream成了轻量级事件总线的最佳人选。业务上,它就是我见过的“一个事件,多方消费”最廉价的实现方式。
5.3 审计日志与近实时时间轴
Stream还有一个容易被忽略的用途:当一条带时间戳的审计流水。用户操作日志、设备上报事件、支付回调状态变更,这类数据的共同特点是:只追加、按时间查、需要保留一段时间。Stream的追加式结构、自动时间戳ID和XRANGE查询能力,几乎是为此量身定做的。
比如你想查“昨天14点到15点之间,用户ID为10001的所有操作”,直接:
XRANGE user:audit:10001 1700000000000-0 1700003600000-0返回结果天然按时间排序,不需要额外索引。配合MAXLEN设置一个合适的保留长度,比如XADD ... MAXLEN 100000,老数据自动被裁剪,内存可控。这种场景如果放在数据库里,你会忍不住加索引、写分页查询,而在Redis Stream里几个命令就完事了。
5.4 用XCLAIM做死信与延迟重试
Stream本身没提供“死信队列”这个高级概念,但你可以用PEL + XCLAIM很便宜地实现类似机制。消费者处理消息失败时,不XACK,消息会一直留在PEL里。你再启动一个独立的“救护车”任务,定时用XPENDING扫出超时未确认的消息,再用XCLAIM把它们转移到一个特定消费者名下,重新投递。
如果消息被重新处理了多次还是失败,你可以再把它XADD进一个专门的死信Stream,比如order:events:dead,然后人工介入。这套方案不依赖任何额外组件,全部在Redis里完成。对于业务团队来说,“有死信可以查”和“没有死信只能翻日志”是完全不同的运维体验。
6. 亲手踩过的坑,和一点点经验
6.1 不要在Stream上设计“严格有序且可并行”的消费
Stream不同消费者组之间是独立并行的,但同一个组内的多个消费者如果同时处理消息,消息的处理完成顺序无法保证。原因很直观:worker-A读到了消息1,worker-B读到了消息2,A处理得慢、B处理得快,最终消息2先完成。如果你的业务要求严格顺序,比如“先扣库存再生成订单”,那同一个组内就不能开多个消费者,或者必须按业务键分片到不同的Stream。
有一种折中做法:同一订单的消息只发给同一个worker。Stream没有Kafka那种partition级别绑定,但你可以手工按订单号散列到多个Stream key,再用一致性哈希把每个key的消费者固定下来。这个方案是可行的,只是需要额外设计。
6.2 别把“读到了”当成“处理成功了”
这是很多新手最容易犯的错。XREADGROUP返回了消息,不代表这条消息就是你的了。Redis的投递语义是“起码一次”,也就是说正常情况下不会丢消息,但在极端场景下可能重复投递:消费者处理完业务逻辑,还没来得及XACK就宕机,重启后PEL里的消息会被重新投递给这个消费者,业务就可能重复执行。
所以Stream消费端的业务逻辑必须做幂等。以订单积分为例,处理前先查一下“这条订单有没有加过积分”,加过了就直接XACK跳过;或者把消息处理结果写进一个幂等表。不要把“不重复”寄托在消息系统上,消息系统只能保证不丢,重复交给业务自己挡。
6.3 阻塞读和BLOCK参数要配合好
XREADGROUP的BLOCK参数是毫秒级阻塞,但Redis的阻塞等待不是“等到有消息就立刻唤醒”,而是基于上一个消息写入事件唤醒等待中的连接,逻辑上存在轻微延迟。实际项目里,如果业务对消费延迟极度敏感,不要只依赖BLOCK长轮询,可以适当加一个短轮询兜底。也不要在一个连接上同时阻塞读多个Stream然后指望等待时间是叠加的,实际上它取的是最大阻塞时间。
还有一个容易踩的细节:BLOCK设为0表示永久阻塞,但如果组里已经很久没有新消息,这个连接会一直挂着,Redis客户端连接池的连接数会悄悄被占满。生产环境我一般设3000到5000毫秒,读不到就返回,自己循环重试,既控制了延迟又避免连接长时间占用。
6.4 内存管理:Stream不是无限日志
Stream默认不会删除历史消息,消息量持续增长时,内存会一路走高。我的建议是写数据时就带上MAXLEN,而不是事后才想起来剪:
# 保留最多10万条,近似裁剪 XADD order:events MAXLEN ~ 100000 * order_id 10001 status paid这里的~是近似裁剪符号,它不会严格删除到精确长度,但性能和内存收益更好。需要特别提醒的是,一旦裁剪,历史消息就真的没了,审计场景要评估一下保留量是否够用。另一个思路是配合多级存储,最近三小时的消息留在Redis Stream里,更老的定期用XRANGE导到数仓,这样两边都舒服。
6.5 XGROUP CREATE的起始位置别搞错
创建消费者组时,你可以选择从Stream头开始读,还是只读新消息:
# 从Stream的第0条开始,能回溯历史 XGROUP CREATE order:events group:all 0 # 只处理创建组之后的新消息 XGROUP CREATE order:events group:new $很多团队在这里想当然地用$,结果接手一条已经在跑的Stream后,发现历史消息全部没消费,排查半天才发现原来是组创建位置选错了。从业务恢复角度来说,如果只是想从断点续跑,用具体消息ID可能更精准。这里建议每次创建组前用XINFO GROUPS确认游标位置。
6.6 关于“多少人能消费”的直觉纠错
刚上手时我总以为Stream支持“一条消息发给所有人”,实际上只说对了一半。不同消费组之间确实都能读到同一条消息,但同一个消费组内部,一条消息只会分给组内一个消费者。也就是说,同一个组里的多个消费者是“竞争关系”,不是“广播关系”。如果你想让多个worker都处理同一批任务,必须建多个消费组,而不是在同一个组里挂更多消费者。很多团队觉得“加消费者就能提升处理速度”,加了之后发现重复消息变多,其实是没搞懂这一层。
7. 关于选型,我现在的判断
结合这几年的使用体感,我对Redis Stream的定位越来越明确:它解决的是一个“质量”问题,而不是“规模”问题。它不会让你的消息系统变快,也不会让Redis变成Kafka,但它能让Redis里的消息传递变得可确认、可回溯、可分组,把此前靠业务团队手工修补的可靠性,下沉成数据结构级别的原生能力。
在实际项目里,凡是消息量在每秒几千级别以内、已经有了Redis基础设施、又对消息丢失零容忍的团队,我都优先建议先试Stream。理由很现实:它不需要额外维护一套中间件,生产端和消费端的改动都很小,遇到问题还能用XRANGE直接在Redis里查消息,排障体验比List队列好太多。
等到哪天你发现单个Stream的吞吐不够、内存成本压不住、或者消息需要跨团队保存几个月,再考虑把核心链路迁到Kafka也不迟。Redis Stream的价值,本来就不是让你抛弃Kafka,而是让你在绝大多数不需要Kafka的场景里,也能拥有一套体面的消息解决方案。