1. 消息到底在哪个环节丢的:先给可靠性画一张路线图
做Kafka的人几乎都听过一句话:“Kafka不丢消息”。可真到了生产环境,因为消息丢失半夜被叫起来排查的人不在少数。为什么会这样?因为“Kafka不丢消息”从来不是一个默认事实,而是一个需要在生产端、Broker端、消费端三个环节分别做对配置、写对代码之后,才能拿到的结果。任何一个环节用了默认值、图省事,丢消息的隐患就已经埋下了。
一条消息从业务进程产生,到消费者真正把业务逻辑执行完,中间要经过三个握手点:生产者把消息发给Broker并等待确认;Broker在副本之间同步数据;消费者拉取消息后提交位移。这三个点分别对应三个核心原则,搞懂它们,后面所有配置和代码就都有了依据:
- 生产端:没有收到Broker的确认,就不能认为消息已经送达。
- Broker端:副本没有同步完成,就不能认为消息已经持久化。
- 消费端:业务逻辑没有处理成功,就不能提交位移。
这三句话听起来简单,但实际落地时每一个环节都有大量细节。比如生产端你设置了acks=all,但没配min.insync.replicas,Broker端照样可能在你只有一个副本存活时给你回“成功”,消息随后就没了;比如消费端你手动提交位移了,但提交的时机放在“处理业务之前”,那处理逻辑一旦抛异常,这条消息就永远丢了。
另外还需要分清一个概念:Kafka在默认配置下提供的是At Least Once语义,也就是“至少一次”,不保证不重复。为了保证“不丢”,我们接受可能出现的重复消息,重复的问题由消费端幂等来解决。那种“既要完全不丢、又要完全不重复、还要性能拉满”的诉求,在分布式系统里是不存在的。先把这一点想清楚,后面很多设计决策就不会纠结了。
2. 生产端拦截:acks、重试与幂等生产者怎么配合
生产端是消息丢失的第一道防线,也是配置项最多、最容易出错的地方。很多人一上来就把acks设成all,以为这样就万事大吉了,其实这只是第一步。
2.1 acks=all和min.insync.replicas是组合拳
acks参数有三个可选值,先看它们各自意味着什么:
| acks取值 | 行为 | 丢消息风险 |
|---|---|---|
| acks=0 | 生产者发完消息不等待任何确认,立即认为发送成功 | 极高,网络抖动、Broker宕机、分区不可用都会无声无息地丢 |
| acks=1 | Leader副本写入本地日志后即返回成功,不等待Follower同步 | 中等,Leader在Follower完成同步前宕机,选主后消息就没了 |
| acks=all | 等待ISR中所有副本都写入成功后才返回 | 低,但还需要配合min.insync.replicas才有实际意义 |
这里最容易被忽略的是:acks=all并不是“等所有副本写完”,而是“等ISR集合里的副本写完”。ISR是保持同步的副本集合,如果某个Follower落后太多,会被踢出ISR。当ISR里只剩Leader一个副本时,acks=all实际上退化成acks=1,消息的安全保障已经没了,但生产者完全感知不到。
所以生产环境必须同时设置min.insync.replicas,它的含义是“接受写入请求的副本最少要有几个在ISR里”。在Broker端的server.properties里加上:
min.insync.replicas=2如果副本数为3,这个值设置为2是比较常见的组合:允许一个副本宕机或落后,但至少还有一个Follower跟着Leader,消息写入时不会出现“只有一个副本”的裸奔状态。当存活副本数小于2时,Broker会拒绝写入请求,生产者的send会返回异常,这时候业务方会收到报错,而不是收到“假装成功”的确认——宁可写入失败,也不能悄悄丢消息。
2.2 幂等生产者:让重试不产生脏数据
生产端的另一个陷阱是重试。Kafka的Producer内置了重试机制,网络抖动、Broker短暂不可用时自动重发消息,这是保证不丢的重要手段。但重试有一个副作用:如果上一次请求实际已经写入了Broker,只是响应超时,生产者重试就会导致同一条消息被写入多次,产生重复数据。
Kafka从0.11版本开始支持幂等生产者,解决的就是这个问题。开启方式很简单:
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);开启后,生产者会给每条消息附带一个递增的序列号,Broker端会对同一个生产者、同一个分区的消息去重,重复写入的请求会被忽略。需要注意,开启幂等生产者时,acks会被自动升级为all,同时max.in.flight.requests.per.connection不能超过5(这个值控制着单个连接上最多有多少个未确认的请求在途,超过5会与幂等机制的序列号校验冲突)。
有了幂等生产者之后,重试就变得安全了:网络超时导致的重发、Leader切换导致的重发,都不会再产生重复消息。这一点在压测环境里很容易验证——故意断网几秒,重启Producer,再对比写入总量,开启幂等前后的差距非常明显。
2.3 发送回调是唯一的“送达凭证”,别只打日志
写Kafka Producer的时候,最常见的坑是在send方法后不关心返回结果。send是异步的,它只是把消息放进了内部的发送缓冲区,真正发出去是后台线程的事。如果你的业务代码是:
producer.send(record);那消息可能在缓冲区内、可能在网络途中、可能已经到达Broker但响应丢了,你完全不知道。正确的做法是注册Callback回调,在回调里检查异常:
producer.send(record, (metadata, exception) -> { if (exception != null) { // 写入失败,需要走补偿流程:记录日志、写入本地失败表、或者降级处理 handleSendFailure(record, exception); } else { // 发送成功,可以记录发送轨迹 recordSendSuccess(metadata.topic(), metadata.partition(), metadata.offset()); } });不要只在回调里打印一行日志就完事。生产环境的经验是:发送失败的消息,必须有一个兜底存储。最简单实用的方案是搞一张本地消息表,发送失败时把原始消息落库,定时任务扫描这张表重新投递。Kafka自带的重试解决的是“临时性失败”,那些Broker持续不可写、消息格式错误、权限问题之类的异常,重试也救不回来,必须靠业务侧的补偿机制兜底。
这里还有一个隐藏参数值得关注:delivery.timeout.ms,它规定了消息从进入Producer缓冲到发送完成(或失败)的总时间上限。默认值一般是120秒,如果你的重试次数配得特别大,但delivery.timeout却很小,那重试可能还没执行完就被超时截断了,结果还是丢消息。所以调优的时候,retries、retry.backoff.ms、delivery.timeout.ms这三个参数要放在一起看,不要单独调某一个。
3. Broker端存住:副本机制、ISR与刷盘的取舍
消息到了Broker,是不是就安全了?还不是。Broker层面要解决的核心问题是:Leader挂了,消息还在不在。
3.1 副本数的意义:Leader不是保险箱
Kafka的一个分区在物理上会有多个副本,其中一个充当Leader,负责处理客户端的读写请求,其余是Follower,只负责从Leader拉取数据保持同步。消息丢失的一个典型场景是:消息写入Leader后,Follower还没来得及同步,Leader突发宕机。此时Controller会从ISR集合里选出新的Leader,而那条还没来得及同步的消息就永远丢失了。
所以Topic的副本数直接影响可靠性。开发环境经常看到有人用默认的1副本创建Topic,这在生产环境等于裸奔。建议Topic创建时统一设置:
bin/kafka-topics.sh --bootstrap-server kafka-1:9092 \ --create --topic order-events \ --partitions 12 --replication-factor 3创建完用describe命令检查每个分区的副本分配情况,确认没有多个副本落在同一台机器上(如果Kafka节点分布在多个机架,最好还能配置rack感知,让副本跨机架分布,避免整机架断电导致全部副本同时失效)。
bin/kafka-topics.sh --describe --bootstrap-server kafka-1:9092 --topic order-events看输出的Leader和Replicas列。当一个副本同步落后时,它的状态会从ISR列表里被移除,此时该分区就处于“降级”状态。正常的describe输出里,每个分区的ISR列表应该和Replicas列表一致,或者只差少数副本。
3.2 关闭unclean选举:宁可短暂不可用,也不让丢消息
Kafka的副本同步有一个特殊情况:Leader宕机后,如果ISR里没有可用的副本了(比如所有在ISR中的副本都挂了),那是不是要选一个落后很多的Follower当新Leader?unclean.leader.election.enable这个参数就是控制这个行为的:
- 设为true:允许从“不同步”的副本中选Leader,服务可用性优先,但必然丢消息。
- 设为false:不允许选不同步的副本,宁可这个分区暂时不可用,也不丢消息。
对这个参数,网上的讨论很多,我的建议非常明确:生产环境必须为false。
unclean.leader.election.enable=false原因很简单:unclean选举是在“可用性”和“一致性”之间做抉择,选了可用性,就要接受消息丢失的后果。而Kafka本身通过ISR机制已经能处理大部分故障场景——ISR里的副本都还活着的情况下,Leader宕机是可以正常选主的,不需要unclean选举介入。真正触发unclean选举的场景,是多个副本同时宕机这种极端情况,这种时候宁可等副本恢复,也不要选一个落后了一万条消息的Follower上台。如果业务真的在意可用性超过一致性,建议从架构层面考虑,而不是在Kafka这里放开口子。
3.3 刷盘策略:Kafka靠副本兜底,不靠强迫落盘
Kafka的消息写入Leader后,实际上先落在操作系统的Page Cache里,由操作系统异步刷到磁盘。很多人第一次知道这件事时会慌:消息还在内存里,机器突然断电不就丢了吗?确实会丢,但Kafka的设计哲学是:单机断电这种极端情况,交给副本机制来兜底,而不是靠每台机器刷盘。
log.flush.interval.messages和log.flush.interval.ms这两个参数控制的是“多长时间/多少条消息强制刷盘一次”。网上有些教程会教你把它们调成1,让每条消息都立刻刷盘。这么做确实能降低单机断电导致的数据丢失风险,但代价是写入性能断崖式下跌,吞吐量可能掉一个数量级。而即便你做了强制刷盘,也没法保证万无一失——存储硬件、文件系统层面的故障依然存在。
正确的思路是接受“依赖Page Cache + 副本机制”的架构设计,用replication.factor=3、min.insync.replicas=2、acks=all这套组合保证数据在多台机器上都有副本,这样任何单台机器断电、磁盘损坏,数据都还在另外两台机器上。刷盘参数保持默认就好,不要乱调。
Broker层面还有一个小细节:controller节点的稳定性。Kafka的Controller负责分区Leader选举、副本分配等元数据操作,Controller所在的Broker如果经常GC停顿或网络抖动,会影响整个集群的副本同步状态。所以集群节点最好不要混跑其他重型应用,JVM的堆内存也要根据分区数量合理设置,避免频繁FullGC。
4. 消费端守住最后一道关:手动提交与幂等消费
很多人认为消息不丢失是生产端和Broker端的事,消费端顶多就是重复消费。这个认知在面试里还能糊弄过去,在真实生产环境里会出大问题——消费端恰恰是最容易“丢消息”的地方,而且丢了还特别难查。关键在于位移提交的时机。
4.1 自动提交是最大的“假成功”
Kafka消费者默认开启自动提交位移:
enable.auto.commit=true auto.commit.interval.ms=5000消费者每5秒会自动把当前拉取到的位移提交给Broker。问题来了:如果你的业务处理耗时较长,比如一条消息要处理2秒,拉取完一批消息后刚处理到一半,自动提交线程把位移提交了(此时提交的是“已经拉取到的最大的位移”,而不是“已经处理完成的位置”),此刻消费者进程重启,Kafka会认为这批消息已经消费完了,从已提交的位移继续消费。那些“拉了但没处理完”的消息,就这样被跳过了——对业务而言,这就是丢消息。
自动提交的本质是“拉取即消费”的假设,它只适合那些“处理逻辑极快、丢了也无所谓”的场景,比如日志采集、指标上报。但凡消息里带的是订单、支付、库存这类有业务含义的数据,必须关掉自动提交:
enable.auto.commit=false4.2 手动提交的正确姿势与重平衡的坑
关掉自动提交之后,手动提交的时机就变得非常重要。最常见的安全做法是“先业务处理,再提交位移”:
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 200); props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 300000); props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 10000); props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 3000); while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> record : records) { // 第一步:处理业务,例如写入订单库、调用下游系统 processMessage(record); } // 第二步:这一批都处理成功了,才提交位移 consumer.commitSync(); }这个模式保证的是:只有处理成功的消息,位移才会往前推进。如果处理中途抛异常,位移不会提交,下次重启会从上次提交的位置重新消费,最多出现重复,绝不会出现丢失。
但这里要注意两个细节。第一,commitSync是同步提交,性能上会有一点损耗,但可靠性最好,建议默认使用。如果追求性能用commitAsync异步提交,一定要在回调里处理失败情况,并且注意:异步提交的完成顺序和调用顺序并不保证一致,进程关闭前必须补一次同步提交,确保最终位移不丢失。
第二,max.poll.interval.ms这个参数。如果你的业务处理时间很长,超过这个阈值(默认5分钟),消费者会被判定为“失联”,触发重平衡,分区会被分配给其他消费者。重平衡本身不会丢消息,但如果你的位移提交逻辑写得不清晰,重平衡期间的状态就很容易搞混。我见过一个案例:消费者用线程池异步处理消息,主线程poll到消息后丢给线程池就立刻提交位移了,结果线程池还没处理完,消费者发生重平衡,那些正在处理中的消息对应的分区被分给另一个实例,处理结果和位移就完全对不上了。这就是典型的“提交过早”。处理逻辑和提交位移必须保持严格的先后顺序,这条线不能破。
还有一个重平衡相关的坑:处理消息时如果调用了外部接口,而这个外部接口偶尔超时,导致处理时间超过了max.poll.interval.ms,消费者就会被踢出消费组。此时最稳妥的方案是先把处理时间优化控制在阈值以内,如果实在做不到,可以适当调大max.poll.interval.ms,或者在处理长耗时任务时采用“先落库、再异步处理、按状态确认”的模式,而不是在poll循环里同步阻塞。
4.3 幂等消费设计:接受重复,消除重复的影响
前面说过,为了保证不丢,我们接受“有可能重复”。那消费端就要有能力处理重复消息。很多业务同学一听到“要写幂等”就头疼,其实幂等消费没有想象中那么复杂,关键是找到业务的天然幂等键。
以订单系统为例,订单号就是天然的幂等键。消费端处理消息时,先查一下订单表,如果这个订单号已经存在且状态为“已处理”,就直接返回;否则才执行正式逻辑。这样即使同一条消息被重复投递,对业务数据也不会产生叠加影响。
更通用的做法是建一张消息消费记录表,主键设为消息的唯一标识(通常是消息的key,如果没有key,可以用topic+partition+offset拼一个):
消费记录表:msg_id(主键)、topic、partition、offset、业务主键、处理状态、处理时间消费流程就变成:
- 根据消息的msg_id查消费记录表,如果已存在“处理成功”的记录,直接提交位移。
- 如果不存在,开启本地事务,同时写入消费记录和业务数据,保证两者要么都成功、要么都回滚。
- 事务提交后再提交位移。
这种方式把“消息消费”和“业务数据写入”绑定在同一个数据库事务里,天然解决了幂等和一致性问题。虽然对数据库有一点额外压力,但从可靠性角度讲,这是很多公司实际在用的成熟方案。Redis的setnx也可以做类似去重,但因为Redis本身可能丢失数据,如果对可靠性要求高,还是建议落库。
5. 如何证明没丢:监控指标、Lag排查与对账机制
配置都做了、代码也写了,接下来的问题是:你怎么知道线上真的没丢?这一节不讲概念,讲实操:用哪些指标、哪些命令、哪些方法去验证消息的可靠性。
5.1 Broker端核心指标:UnderReplicatedPartitions不能一直大于0
Broker端最关键的一个指标是UnderReplicatedPartitions,它表示“正在同步中的副本数低于正常值”的分区数量。正常情况下这个值应该为0。如果你在监控面板上看到它持续大于0,说明有分区副本同步落后或副本已经不在ISR里,这是消息丢失的前兆——因为一旦此时Leader宕机,那些没同步完的消息就没了。
查看方式有两种。Kafka自带命令行:
# 查看某个Topic所有分区的ISR状态 bin/kafka-topics.sh --describe --bootstrap-server kafka-1:9092 --topic order-events以及通过JMX监控。生产环境建议用Prometheus + Grafana收集Kafka的JMX指标,重点盯这几个:
- UnderReplicatedPartitions:长期大于0要告警。
- IsrShrinksPerSec / IsrExpandsPerSec:ISR频繁收缩扩张,说明有副本经常掉队。
- OfflinePartitionsCount:分区Leader离线数,大于0说明有分区不可用。
- RequestHandlerAvgIdlePercent:请求处理线程的平均空闲率,太低说明Broker负载过高。
- NetworkProcessorAvgIdlePercent:网络线程空闲率,类似。
如果你不想从零搭监控,也可以先用AKHQ这类Kafka可视化工具快速看一下集群状态。AKHQ能直观展示Topic列表、分区ISR状态、消费组Lag信息,还能查看Connector任务状态,排查问题时比命令行高效很多。不过可视化工具适合快速查看,长期监控还是得落到指标系统里。
5.2 消费Lag排查:不是所有高Lag都是消费慢
消费Lag,即消费组当前消费到的位置和生产端最新写入位置之间的差值,是判断消费端是否正常的最直观指标。
bin/kafka-consumer-groups.sh --bootstrap-server kafka-1:9092 \ --group order-service --describe输出里看LAG列,代表这个消费组还有多少条消息没消费完。Lag持续增长,说明消费速度跟不上生产速度,正常情况下应该趋近于一个较小值。
这里要区分两种Lag高的原因。一种是消费组整体处理能力不足,表现为所有分区Lag都在涨,这时候加消费者实例数有可能解决(前提是分区数大于消费者数)。另一种是某个特定分区的Lag特别高,其他分区正常,这大概率是“消息卡在某个业务处理上了”——比如消费端处理某条消息时抛了异常,或者走到了死信分支,但位移一直没提交,导致这个分区后续消息都堵住。排查这种问题,先看消费组日志里有没有异常堆栈,再看是不是有消息处理超时,不要一上来就加实例。
还有一个容易被忽略的点:生产端和消费端的时间戳。如果生产者的消息时间戳和Broker服务器时间不一致,通过时间维度去看Lag会产生误导。排查时尽量以“位移差值”为准,不要以时间判断消费进度。
5.3 对账表:把“感觉没丢”变成“数据证明没丢”
监控指标能告诉你“损坏风险”,但有些消息丢失是“静默”的——配置错、代码逻辑错、位移提交时机错,它不会产生任何告警,只有业务数据对不上时才会被察觉到。这时候就需要对账机制兜底。
我的做法是:在生产端写消息时,同时往一张消息轨迹表里插一条记录,记录这条消息的topic、partition、offset、业务主键、发送状态。消费端消费成功后在轨迹表里更新状态为“已消费”。然后定期(比如每小时)跑一个对账任务,扫描那些“已发送但长时间未消费”的消息,人工确认是否真的丢了,再决定是否需要重放。
这套东西做起来不复杂,但对业务数据的一致性非常有用。它解决的不只是“消息丢失”,还包括“消息到达了但业务没处理”的情况。很多线上事故的定位,最后都是靠这样的对账日志还原出问题的全貌。如果你的业务体量不大,可以不用建完整的对账平台,但至少要保证生产端和消费端都有落地的日志表,为排查留一条路。
6. 几个容易让人迷糊的问题:事务、Exactly Once与Kafka和RabbitMQ的差异
关于Kafka消息不丢失,还有一些高频出现的疑问,这里集中梳理一下。这些问题在面试里也经常被追问。
6.1 Kafka的事务能解决“不丢不重”吗?
Kafka从0.11起支持事务API,通过transactional.id配合事务消息,可以实现跨分区原子写入。但很多人的理解是“开启了事务,消息就不会丢也不会重复”。这个理解是有偏差的。
Kafka事务解决的是“生产端写入多个分区的原子性”问题,以及“Consume-Transform-Produce”这种流处理场景下,消费位移和写入结果之间的原子性问题。如果你只是普通的“业务系统 -> Kafka -> 业务系统”这种消息中转,开启事务并不会让消息可靠性变得更高,反而会引入大量性能开销(事务协调器的提交确认、事务日志的刷盘等)。
真正能让“生产到消费”整体具备Exactly Once语义的场景,是Kafka Streams这类流处理框架:它把“读取Kafka消息、计算结果、写回Kafka、提交位移”这四步包在一个事务里,要么全成功,要么全回滚。普通消费者自己写业务逻辑时,要么用“事务+幂等表”,要么接受“至少一次+业务幂等”,不必非要用Kafka的事务API,复杂度不成比例。
6.2 Kafka和RabbitMQ在消息可靠性设计上有什么本质区别
这个对比经常出现在面试题里,但日常工作中理解它也有实际价值。两者在可靠性设计上的核心差异可以概括为:
- Kafka:以日志为核心,消息有序存储在分区里,消费者通过位移自主控制消费进度,可靠性建立在“多副本 + 位移提交时机”上。它的设计假设是消费者可能会离线很久,Broker不会为了消费者删除消息,所以天然适合“消息延迟消费”“消息回溯”等场景。
- RabbitMQ:以队列为核心,消息投递给消费者后,如果消费者未确认,消息不会被删除(unacked状态)。它的可靠性更多依赖“手动ack机制”和“持久化交换机/队列/消息”三件套。
所以如果你要用RabbitMQ保证不丢,必须同时开启持久化,而且消费者一定要手动ack;Kafka这边则不强制消费者立即确认,重点反而在“位移提交不要早于业务完成”。同样是“消息不丢失”,两边的设计思路和关键参数完全不同,这也是为什么网上总有教程强调“不要用Kafka的自动提交,就像不要用RabbitMQ的自动ack”是一样的道理。
6.3 消费端遇到“毒丸消息”怎么办
所谓毒丸消息,就是单条消息本身内容有问题,导致消费端每次处理都抛异常。如果不处理,位移提交不了,消费线程会一直卡在这条消息上,整个分区都被堵死。这种场景不算“消息丢失”,但它的危害和消息丢失一样大——后面的消息全部积压,Lag直线上升。
常规解法是把这类消息单独隔离:捕获异常后,判断是否是业务异常,如果是,将消息内容落库(死信表),然后手动提交位移,继续消费后面的消息。但这有个前提:你确认这条消息确实“必定处理不了”,否则会掩盖真正的逻辑Bug。经验是:第一次遇到异常时先重试几次,重试仍失败的再进死信表,死信表要有专门的告警,不能让消息悄悄进去就完事了。
我在实际项目中就是把“保存死信 + 更新消费记录 + 提交位移”串在同一个流程里,等于是给消费链路加了一个旁路,既不阻塞主流程,又不放过问题消息。线上跑下来,比一直卡住整个分区要省心得多。
最后再分享一个经验
做了几年Kafka相关的系统,经历过半夜被叫起来查消息丢没丢的日子,我的体会是:消息可靠性这件事,最怕的是“以为做好了”。acks=all配了,min.insync.replicas没配,等于白配;手动提交写了,但放在处理逻辑前面,等于没写;监控面板搭了,但只看CPU和内存,关键的分区ISR状态一个都没盯,等于白搭。
给新团队的建议是:先把本文提到的三个环节的配置和代码都过一遍,然后在压测环境做一次破坏性演练——杀掉一个Broker节点,重启消费端服务,批量灌数据,看看有没有消息丢。演练通过了,再提“不丢消息”这回事也不迟。真正稳的Kafka使用方,不是靠某一次配置调优,而是靠一套“配置、代码、监控、对账”四位一体的机制在兜底。