消息队列这四个字,在很多团队眼里就是个“发件箱”。订单创建成功了,往队列里丢一条消息,库存、积分、短信各取所需,谁有空谁来消费。这个理解大方向没错,但如果你真把它当成一个普通发件箱来用,生产环境迟早教你做人。
做了快十年的消息中间件运维和架构设计,我最大的感受是:消息队列这个“产品”本身并不复杂,复杂的是你把它放在什么样的业务链路里,以及你有没有提前想清楚丢消息、重复消息、顺序错乱这些极端情况。这篇文章不准备给你背一遍Kafka、RabbitMQ、RocketMQ的官网参数,而是想从一个真正落地过、也踩过坑的从业者角度,把消息队列的选型、原理、重复消费这些最容易被忽视却又最致命的问题,一次性讲透。
面向的读者,我默认你是后端开发、系统架构师或者负责中间件运维的同事。如果你是刚接触消息队列的学生,这篇文章也能当一份比较完整的入门图谱来读,但请做好反复回看的准备,因为里面不少结论都是生产环境用血泪换来的。
1. 业务系统为什么非要有消息队列
1.1 把串行请求改成并行异步,系统的韧性完全不一样
先抛一个最简单的场景。电商下单,用户点击“提交订单”之后,后端要干的活包括:扣库存、生成订单、发积分、发短信、更新报表数据。如果全部同步调用,假设每个环节50毫秒,五个环节串起来就是250毫秒,其中任何一个下游系统抖动,用户的请求就要跟着一起卡住。更要命的是,像短信、积分这类系统和核心交易链路并无强依赖——短信通道故障,难道还不让用户下单了?
引入消息队列之后,下单主链路只做两件事:写订单数据、往队列里放一条“订单已创建”的消息。库存、积分、短信这些下游业务自己去订阅消息,能处理就处理,处理不给力就先积压着。用户响应时间从250毫秒降到100毫秒以内,整个系统的可用性也不再被下游的短板卡死。这个思路就是我们常说的“异步解耦”,它是消息队列在企业系统里最基础、也最持久的使用方式。
1.2 流量削峰填谷,相当于给洪峰修了一个水库
秒杀、大促、预约抢购这类场景,流量特征永远是“尖峰突刺”。你可能平时每秒只有1000个请求,到了开场那一秒暴涨到10万。如果让下游数据库直接扛这10万QPS,别说MySQL,再贵的机器也得跪。与其硬扛,不如把瞬间的写入请求全部先“吞”进消息队列里,让后端系统按照自己能够承受的速率缓慢消费。
这就是削峰填谷。你会发现,消息队列在这里的角色更像一个水库:上游的洪流涌进来,它在中间拦一道,下游的渠道按照恒定流量往外放水。这样做的好处非常明显,你不用为了一年只有几次的峰值流量去把整条链路扩容到十倍规模,只要保证队列的写入吞吐和Broker的存储容量足够即可,成本省下了,系统也稳了。
1.3 数据分发与系统解耦,是微服务架构的地基之一
老一代的单体系统里,A系统要一份数据,直接连数据库查B系统的表,或者B系统提供一个HTTP接口。这种方式的隐患在于,下游每多一个需求方,上游就要改一次代码;下游数据库结构一变,上游直接崩掉。消息队列天然是“发布-订阅”模型——生产者只管把消息发到Topic里,不关心谁在听;消费者按需订阅自己关心的消息类型,不关心消息是谁发的。这种模式让系统的边界变得非常干净,两个团队之间只需要约定好消息格式,其他一切互不打扰。
我自己参与过的一个供应链项目就是典型例子。订单系统产生的状态变更消息,被仓储系统用来决定是否发货、被财务系统用来记账、被客服系统用来展示物流进度,还同时被数据团队拉到数仓里做分析。四个消费方,一个Topic,订单系统每次发布消息,根本不用知道外面多了谁少了谁。
当然,你也得清醒一点:消息队列不是万能的银弹。它引入了“最终一致性”的复杂度,消息丢失、重复消费、顺序错乱都是它的固有特点,不是产品有Bug。这就是为什么很多人调侃“没有经历过消息积压的架构师,不足以谈人生”。
2. 消息队列产品介绍:拆开内核看概念,不装懂
2.1 从发件箱到邮局体系:消息队列里的角色划分
为了把概念讲清楚,我喜欢用现实中的邮局来类比。消息队列由三个核心角色组成:
- 生产者(Producer):写邮件的人,只负责把信投进邮筒,不关心收件人具体怎么处理。
- Broker:就是邮局本身。它接收消息、存储消息、按规则把消息投递给消费者,是整个系统的中枢。
- 消费者(Consumer):取信的人。消费者根据自己的订阅关系从Broker拉取消息,然后执行具体业务逻辑。
这里要特别强调Broker,因为很多初学者会把消息队列理解为“一个进程”,实际上生产环境里的Broker通常是一个集群。以Kafka为例,一个Kafka集群由多台Broker节点组成,通过ZooKeeper或KRaft协议协调元数据,单节点挂掉后集群仍然可以对外提供服务。这种多节点结构才是它能够扛住高并发、高可用的根本原因。
2.2 Topic、Partition与消费组:并行度从哪来
Topic是消息按照业务逻辑分类的大类,比如“订单消息”、“支付消息”。但光有Topic还不够,Kafka和RocketMQ这些现代消息系统内部,会把一个Topic切分成多个Partition(分区),每条消息根据Key的哈希或轮询策略,被写入到某个分区中。
为什么要做分区?这里有两个决定性原因。第一,并行写入:如果所有消息都写在一个文件里,写入瓶颈很快就会卡住;分区之后,不同分区可以落在不同磁盘、由不同线程读写,写入吞吐成倍提升。第二,并行消费:Kafka保证的是“分区内有序”,只要一个分区被同一个消费组内的消费者实例拿到,在这个分区内消息就是按顺序被处理的。所以,如果你希望某个业务的消息严格有序,必须把相同业务Key(比如同一个订单号)的消息路由到同一个分区去。
消费组也是一个绕不开的概念。同一个组内的消费者实例会分工消费同一个Topic下的分区,一条消息只会被组内某一个实例处理。不同消费组之间则互不影响,都能拿到全量消息。这就是“集群消费”和“广播消费”的底层逻辑。一个Topic有12个分区,挂着4个消费者,那每个消费者平均处理3个分区;如果你只开了1个消费者,那12个分区的活就得它一个人干,消费能力自然跟不上。
2.3 Offset与ACK机制:进度到底谁来管
Offset(偏移量)是消费者在某个分区中的消费位置。消费者每处理完一条消息,就要向Broker提交一次Offset,告诉Broker“这个位置之后的消息还没处理,下次从这继续”。而ACK(Acknowledgment)是消费者对消息处理结果的确认回执。
如果你的代码里设置了自动提交位移,也就是enable.auto.commit=true,消费者拉取到消息之后就会定期提交Offset,不管这条消息的业务逻辑是否执行成功。这是最简单的配置,同时也是最常见的消息丢失来源。生产环境我强烈建议你关闭自动提交,改为手动提交位移,并且要在业务逻辑成功处理完之后再提交。虽然手动提交会让代码复杂一些,但这是唯一能让你把“消息可靠性”握在自己手里的方式。
2.4 存储与可靠性的底层逻辑,决定了产品的性格
不同消息队列的存储设计,直接决定了它们的性格差异。
Kafka存储使用的是追加式的分段日志,每个分区对应一个目录,内部包含多个Segment文件。消息写入时只做顺序追加,不修改旧数据,再辅以操作系统的Page Cache加速读写,所以它的吞吐量在所有消息队列里是一骑绝尘的。但代价是它的消息在写入完成后不会立即对消费者可见,而且要等消息被消费后过一段时间才会被清理。
RabbitMQ则更像传统的消息代理,消息可以配置为持久化,但它在设计上是面向低延迟、内存优先的。你没看错,RabbitMQ擅长的是灵活的路由策略和毫秒级低延迟,但当消息大量积压时,它远没有Kafka那种硬扛上亿条消息堆积的能力。
RocketMQ吸收了Kafka的分区思想,但存储上做了一些折中:它把每个Topic的消息连续写入一个CommitLog,然后通过ConsumeQueue来逻辑索引。RocketMQ的好处是消息堆积能力很强,同时也支持事务消息、延迟消息这些高级特性。Pulsar则比较特别,它把存储层单独抽出来,用BookKeeper做持久化存储,实现了存储与计算的完全分离,多租户和跨地域复制能力很强,但整体架构复杂度也会更高。
理解了这些底层机制,你才能理解为什么Kafka这么能扛积压,为什么RabbitMQ不适合做离线大批量消费,为什么选型的答案不能只看网上的一张对比表,还要结合你自己的业务特点。
3. 四大热门消息队列选型实战对比
3.1 Kafka、RabbitMQ、RocketMQ、Pulsar核心参数横评
选型之前,先看一张我整理的关键特性对比表。这张表的信息来自我在多个生产项目的实测体验,以及社区公开的性能测试报告,适合作为你选型的第一参考:
| 对比维度 | Kafka | RabbitMQ | RocketMQ | Pulsar |
|---|---|---|---|---|
| 开发语言 | Scala/Java | Erlang | Java | Java |
| 消息模型 | Topic/Partition | Exchange/Queue | Topic/Queue | Topic/Subscription |
| 典型吞吐量 | 极高,百万级TPS | 中等,万级TPS | 高,十万级TPS | 高,十万级TPS |
| 消息延迟 | 毫秒~秒级 | 微秒~毫秒级 | 毫秒级 | 毫秒级 |
| 堆积能力 | 极强,TB级轻松 | 较弱,积压后会丢性能 | 强,亿级消息没问题 | 极强,存储与计算分离 |
| 顺序消息 | 分区内有序 | 单队列有序 | 队列内有序 | 单订阅有序 |
| 定时/延迟消息 | 不支持原生 | 插件/延迟队列 | 原生支持18个延迟级别 | 原生支持 |
| 事务消息 | 支持(需配合幂等) | 支持(事务消息) | 支持(半消息机制) | 支持 |
| 消息轨迹 | 需额外实现 | 需插件 | 原生支持 | 原生支持 |
| 流式计算 | 强,配合Kafka Streams/Flink | 弱 | 一般 | 强 |
| 运维复杂度 | 中,需要管理分区 | 低 | 中 | 高 |
你会发现,没有完美的产品,每个队列都有它的主场和短板。Kafka家底厚实但组件繁多,RabbitMQ灵活轻量但一遇积压就露怯,RocketMQ功能全面但是Java技术栈的人用起来才顺手,Pulsar理念先进但也意味着团队要有相当强的掌控力。
3.2 不同业务场景下怎么选,我给你一套判断顺序
网上聊消息队列选型,动不动就是一张巨大的表格,看完更晕。我自己在项目里习惯按“业务诉求优先级”来一步步过滤,操作性强很多:
- 如果你的主要诉求是大数据管道、日志收集、流量削峰、流式计算,那么直接选Kafka。它设计目标就是高吞吐、持久化、顺序追加,你在这些场景下基本不会踩坑。别拿它去做RPC级别的延迟敏感业务,它的延迟在超高吞吐下并不稳定。
- 如果你的诉求是传统业务解耦、异步通知、小而灵活的消息路由,且预计消息量不大,RabbitMQ是正确的选择。它的管理界面友好,路由模型灵活,团队成员上手成本低。但要记住一个铁律:RabbitMQ不要攒消息,积压超过一定量级,性能会断崖式下跌。
- 如果你在Java技术栈体系内,业务场景又比较复杂——既要削峰、又要延迟消息、又要事务消息、还要消息轨迹,那RocketMQ几乎是为你量身定制的。阿里内部大规模应用验证过,社区的生态也很成熟。唯一需要注意的是,RocketMQ对客户端版本比较敏感,升级时要先做好兼容性测试。
- 如果你的场景是多租户、跨地域复制、消息积压常态化,预算和运维团队都不差,那可以认真考虑Pulsar。它的架构非常先进,但也因为先进,踩坑资料少、招人难,小团队一定要慎重评估自己的掌控能力。
3.3 选型避坑:这三个错误我见得太多了
先说第一个坑:用Kafka做业务系统的默认队列。Kafka吞吐高不代表它适合所有业务。订单、支付这种核心交易链路的消息,需要的是低延迟、可管理、能回溯,Kafka在这块并不是体验最好的那个人。有些团队把Kafka架起来,结果查询单条消息要翻半天日志,消息轨迹全靠自己实现,维护成本远比想象中高。
第二个坑:迷信RabbitMQ的“轻量”标签,把它用于核心数据的中转。RabbitMQ确实轻量,但轻量意味着内存管理和磁盘缓冲能力有限。一旦流量突发堆积,RabbitMQ会迅速进入内存高水位报警,甚至触发流控,把消息吐回生产者。这不是产品不好,是使用场景没匹配上。
第三个坑:忽视Pulsar和Kafka的元数据集群。很多人看到Pulsar的存算分离,就以为Broker是无状态的可以随便扩容,实际上它对元数据服务BookKeeper集群的要求极高,BookKeeper节点抖动会引发大面积写入失败。Kafka则在新版本里用KRaft协议替代了ZooKeeper,老经验不完全适用,升级时务必看官方迁移文档。
4. 消息队列重复消费问题:躲不掉的分布式宿命
4.1 重复消费为什么是常态而不是意外
很多业务团队第一次遇到重复消费,第一反应是“消息队列出Bug了”。其实恰恰相反,在分布式环境下,重复消费几乎是必然会发生的事件。原因可以归纳成三类:
- 消息已处理,但Offset没来得及提交。消费者拉取到消息,执行完业务逻辑,正要提交Offset时进程崩溃了。重启后,Broker认为自己没投递成功,会重新把这条消息发给消费者。这是最典型的重复消费场景。
- 网络超时导致的重试。消费者处理消息耗时较长,Broker长时间没有收到ACK确认,就会触发重试投递。消息本身可能已经被处理完了,但确认信号丢了,于是又投递一次。
- 消费实例变更。在Kafka里一个消费者挂掉或者新增消费者加入消费组,会触发Rebalance,分区重新分配。分配的那一刻,有些消息可能已经被上一个实例拉取但未处理完,新实例就会重新拉取一遍。
所以,不要抱着“我刚好不会遇到”的侥幸心理。你应该默认“每条消息至少会被投递一次”,然后把系统设计成“即使重复处理也不会出问题”,这就是所谓的消费幂等。
4.2 生产端如何配合:幂等生产者和唯一消息ID
重复消费的问题,源头其实可以追溯到生产端。如果生产者由于网络抖动,发送消息时发生超时重试,Broker上就可能被写入两条一模一样的消息。Kafka本身提供幂等生产者的能力,你只要在生产者配置里设置enable.idempotence=true,Kafka就会通过生产者ID和序列号机制,自动过滤掉由重试引起的重复消息。
不过这只是解决了“Broker存储层”的重复,跨系统的重复你是挡不住的。更可靠也更简单的方法是:在生产消息时携带一个全局唯一的业务ID,比如订单号加消息类型组合而成。下游消费者拿到这条消息后,先检查这个ID是否已经处理过,处理过就直接忽略。这个方案虽然要多写几行代码,但它能覆盖从生产到消费的全链路重复场景,是性价比最高的手段。
4.3 消费端幂等:三种最实用的设计模式
消费端幂等,核心思想是“无论消息重复多少次,最终的数据状态是一致的”。我在生产项目里经常采用的是下面三种方式:
第一种是数据库唯一键去重。在处理消息时,往一张“消息处理记录表”里插入一条记录,把消息ID作为唯一主键或者唯一索引。如果这条ID已经存在,插入就会失败,业务代码捕获冲突后直接标记消息成功。这个方案严格有效,代价是多一次数据库写操作。
第二种是状态机校验。很多业务的更新操作并不是单纯的累加,而是有明确的状态流转,比如订单从“已支付”到“已发货”。在消费逻辑里,先查询当前业务对象的状态,如果已经处于目标状态或更终状态,就直接跳过。这种方式非常适合订单、物流这类状态明确的业务,几乎不增加额外成本。
第三种是基于Redis的幂等锁。在消费前先执行SETNX(Key为消息ID,Value为时间戳,过期时间合理设置),如果返回成功,说明这条消息是第一次处理,继续向下执行;如果不是,说明已经处理过了,直接返回。这个方案性能好,但有前提——Redis本身必须是大致可靠的,否则你判断幂等的依据本身就可能出错。
这三种方式没有绝对的优劣,关键看你们的业务形态和已有技术栈。如果让我给一个推荐优先级:能用数据库唯一键就用数据库唯一键,这是最简单也最不容易出错的方案;状态机校验适合天然存在状态流转的场景,往往可以和业务逻辑合二为一。
4.4 手动提交Offset的正确姿势,细节别搞错
重复消费的直接诱发点就在Offset的提交策略上。生产环境我一律建议关闭自动提交,代码上按下述步骤操作:
- 消费者拉取一批消息,比如100条。
- 逐条(或按小批量)执行业务逻辑。
- 每条消息处理成功之后,更新本地消费进度。
- 这批消息全部处理成功后,再一次性向Broker提交这批消息中最后一条的Offset。
这里最关键的是,提交Offset必须严格晚于业务数据处理成功之后。有些初级开发图省事,拉取完消息立刻提交Offset,再慢慢处理业务逻辑。这会造成什么后果?如果业务逻辑处理到一半进程重启了,Broker认为这些消息已经消费过了,不会再投递,那这部分业务数据就悄悄丢了。所以,请牢牢记住一句话:宁可重复消费也不要丢消息,重复消费可以用幂等来兜底,丢消息却等于数据事故。
5. 消息队列落地实战与调优经验
5.1 Kafka生产者和消费者调优,这几个参数先记住
很多团队部署Kafka,默认配置一跑就是半年,直到出问题才想起来调优。我建议至少先把下面这组生产者和消费者参数吃透:
对于生产者,核心配置是acks=all、retries(重试次数,建议大于3)、enable.idempotence=true、linger.ms(等待时间,建议5到50毫秒)、batch.size(批量大小,16KB起)以及compression.type(建议LZ4或ZSTD)。acks=all保证消息写入所有ISR副本后才算成功,这是不丢消息的底线。linger.ms很多人以为设得越小越好,实际不然——它允许生产者把多条消息打包成一批再发送,稍微的延迟可以换来吞吐量的大幅提升。如果你的业务对延迟要求在100毫秒内,设个10到20毫秒是完全安全的。
对于消费者,核心参数是enable.auto.commit=false、max.poll.interval.ms(两次poll最大间隔,默认300秒)、max.poll.records(单次poll返回的最大记录数,建议500以内)和session.timeout.ms。这里要特别提醒:max.poll.records设得太大,消费者拉取了一大堆消息,但处理速度跟不上,两次poll的间隔一旦超过max.poll.interval.ms,就会被判定为消费异常,触发Rebalance。这种问题在生产环境非常隐蔽——表面上你的消费者进程活着,实际上它已经被踢出消费组了。
5.2 RocketMQ的高阶特性:延迟消息和事务消息的坑
RocketMQ之所以在很多业务系统里受欢迎,延迟消息功不可没。它原生支持18个延迟级别,比如1秒、5秒、10秒、30秒、1分钟、2分钟……直到2小时。下单后30分钟未支付自动取消,这个需求用RocketMQ的延迟消息实现,就是发送一条“延迟30分钟被消费”的消息,定时触发检查订单状态。
但这里有一个隐蔽的坑:RocketMQ默认的延迟级别粒度是固定的,你没法任意指定“45分钟后执行”,只能选择预定义级别。如果要更强的定时能力,需要自己实现时间轮或借助外部调度框架。事务消息方面,RocketMQ的半消息机制很强大——生产端先发送半消息,Broker存储但不投递,然后执行本地事务,根据事务结果提交或回滚半消息。但要注意,如果你没有实现可靠的回查接口,事务消息在极端情况下可能会卡在中间状态。所以事务逻辑一定要做好幂等设计和状态持久化。
5.3 运维监控体系:这四个指标决定了队列是否健康
消息队列平时很安静,但一旦出问题就是大问题。我自己的监控基准线是这四项:
- 消息积压量:消费位点与最新位点的差值。积压量持续上涨,说明消费能力已跟不上生产速度,需要扩容消费者或调查消费耗时。
- 消费耗时P99:单条消息从被拉取到处理完成的延迟分布。P99突然变大,通常意味着数据库或下游接口出现瓶颈,这时候盲加消费者也救不了。
- 消费失败与重试次数:重试次数飙升,往往对应下游服务不稳定或消息格式不兼容。
- Rebalance频率:Kafka消费组的Rebalance频率过高,最常见的原因是消费者处理速度过慢、心跳超时或消费者实例频繁上下线,这种问题会进一步放大消费延迟,形成恶性循环。
监控体系搭建好之后,要做到“积压告警当天响应,消费耗时告警30分钟内定位”,消息队列的运维才不会沦为救火队。
6. 消息队列常见问题与排查技巧实录
6.1 消息丢了、重复了、积压了、乱序了,一张表看全
| 问题现象 | 典型原因 | 排查思路 | 解决方案 |
|---|---|---|---|
| 消息丢失 | 生产端发送失败未重试、Broker刷盘策略过松、消费端自动提交Offset | 查生产日志有没有发送失败、排查Broker刷盘配置、查看消费端是否开启自动提交 | 设置acks=all、开启幂等生产者、关闭自动提交、业务成功后手动提交Offset |
| 消息重复 | 网络超时重试、Offset未提交、Rebalance导致重新拉取 | 查看消费端是否有幂等校验、检查提交Offset时机 | 消费端做幂等(唯一键/状态机/Redis锁),生产端开启幂等,尽量手动提交 |
| 消息积压 | 消费吞吐跟不上、下游接口缓慢、单分区无法并行 | 看消费耗时、看积压曲线、看分区数量与消费者数量 | 优化业务处理逻辑、增加分区与消费者实例、优先保障核心链路消息,必要时先清空过期消息 |
| 消息乱序 | 同一业务消息进入不同分区、失败重试导致乱序 | 查看是否按业务Key路由分区、重试机制是否合理 | 确保同一业务Key发往同一个分区;顺序性要求高的场景慎重开启重试,或由消费端做排序和状态判断 |
6.2 排障实录一:Kafka消费组频繁Rebalance,问题出在一条慢SQL
有一次线上Kafka消费者频繁告警,消费组每隔几分钟就Rebalance一次。表面上看是消费者实例不稳定,但排查了很久,进程和网络完全正常。真正的原因,是消费逻辑里执行的一条数据库查询语句没有加索引,单条消息处理耗时从10毫秒飙到了3秒。消费者处理不过来,poll间隔超时,服务端判定它“失联”,把它踢出消费组,于是Rebalance,分区分给了其他消费者,结果其他消费者一样变慢,再次被踢,陷入死循环。
这个案例告诉我们,Kafka的Rebalance问题并不总是自身配置问题,消费端的下游依赖往往是罪魁祸首。遇到Rebalance,第一件事先拉出消费耗时和慢查询日志,而不是盲目调大session.timeout.ms。
6.3 排障实录二:RocketMQ消息偶发丢失,最后查到是起了多个消费实例
另一个项目里,RocketMQ的消息偶尔会丢失,毫无规律。查遍了生产者和Broker的配置,都没发现问题。最后开了消费组的监控,才发现同一个消费组被两个不同的应用进程注册了,其中一个还是老版本的代码。由于同名消费组共享消费位点,A进程消费了一批消息并提交了Offset,B进程还在按旧逻辑处理同一批消息,处理失败之后又因为位点已经被提交了,无法重试,消息就这么“丢”了。
所以我想强调一点:消息队列的消费组名称必须全局唯一,而且一个消费组的代码最好只部署在一个应用里。别小看这点,它是很多人排查半天才发现不了的隐性坑。
6.4 排障实录三:RabbitMQ一积压就崩,换RocketMQ后稳如老狗
有个做电商供应链的朋友,早期用RabbitMQ做订单消息中转。平时流量小,一切正常,一到促销活动,消息量翻三五倍,RabbitMQ的内存水位就往上飙,紧接着进入流控状态,生产端不断报错。后来他们把核心订单链路切换到RocketMQ,同样是积压,RocketMQ能稳稳地把消息写到磁盘,高峰期过后慢慢消费完。这个案例不是说RabbitMQ不行,而是你要想清楚它适合什么场景——它是低延迟灵活路由的专家,不是海量消息堆积的好手。
7. 我个人坚持的几个消息队列落地底线
看到这里,你可能会觉得消息队列的水很深。确实,它在一个业务系统里看起来只是不起眼的中间件,一旦出问题,影响范围却是全链路级别的。我做了这些年消息中间件,最终沉淀下来的原则其实只有三条。
第一条,所有核心消息必须开启幂等。不管是Kafka的幂等生产者,还是消费端的业务幂等,这笔代码一定不能省。不要觉得“我们业务流程很简单不会重复”,我之前经历过的那点重复消费案例,十个里有九个也是这么想的。
第二条,消息可靠性宁可牺牲一点吞吐也不能丢数据。生产端acks=all加上消费者手动提交Offset,这是底线。除非你做的只是日志采集这种允许丢失的场景,否则一秒钟都不要放弃这条底线。
第三条,上线前一定要做故障演练。把Broker节点故意停掉一台、把消费进程故意杀掉、把消息积压量人为拉高,测试系统在这三类故障下能不能自动恢复、会不会产生数据错乱。很多团队部署消息队列没出过事,不是系统设计得好,只是运气好没有触发边界条件。
最后再分享一个小技巧。消息队列的位点信息要定期导出备份。有一次线上Broker的磁盘故障导致日志损坏,位点信息全部丢失,靠备份才把消费进度拣了回来,避免了一次完整的数据重放。这个动作成本极低,关键时刻能救命。