
消费端保证消息不丢失Kafka里最容易被误解的一个话题。很多文章一上来就让你把enable.auto.commit设成 false然后手动提交位移说得好像只要这么做了消息就一定能安稳落地。但我在生产环境里排查过太多消费端异常结论是手动提交只是入场券真正决定消息丢不丢的是你如何处理业务结果、位移提交、重试和重新平衡rebalance之间的关系。这篇文章想把消费端“不丢失”这件事从机制到实践彻底拆开讲清楚丢消息的根源、位移提交的坑、幂等的必要性以及一套我亲测有效的两阶段消费方案。这套方案里的本地待处理表加外部位点设计在不少讲Kafka的文章里是真没见过但它恰恰是高可靠消费的核心。1. 先把问题说透消费端丢消息到底丢在哪1.1 我见过最多的“丢消息”其实是这三种先说结论Kafka消费端在默认机制下是at-least-once至少一次不是exactly-once精确一次。也就是说从Kafka自己的语义看消息几乎不会被“机制性丢失”真正的丢多数发生在业务处理环节。我梳理了自己这些年遇到的消费端丢消息故障基本逃不出下面三种第一种消费者在处理消息时抛了异常代码里直接 catch 住然后接着向后消费这条消息从业务视角就是丢了Kafka却毫不知情它以为你处理完了。第二种位移提交时机错了。比如你先提交了位移再去写下游数据库结果下游写库失败等消费者重启后Kafka认为这条消息已经消费过直接从下一条开始读那这条失败的数据就永久丢了。第三种消费者处理太慢触发 rebalance 被踢出消费组。被踢出后你之前拉取但还没提交位移的消息会被分配给其他消费者重新消费而你自己处理到一半的业务结果可能已经写入系统于是造成重复处理。把场景放大一点如果另一个消费者处理的结果和你这边不一致业务上就会出现数据错乱甚至覆盖掉正确数据。这三种场景里第一种和第三种严格说是“漏处理”或“重复处理”第二种才是真正意义上的“丢”。但在业务视角只要系统没有稳定地处理完所有消息都算不保证不丢失。所以后面所有方案本质上都是围绕这三个场景做文章。1.2 为什么说消费端默认能力只是“至少一次”Kafka消费端的核心状态是一个叫 offset位移的东西它记录的是消费者组在某个分区上已经“确认处理完”的位置。主题里的消息在分区内按 offset 顺序排列消费者从头往后读读完一条位移就往前走一步。问题在于Kafka只负责推进位移它完全不知道你的业务是否真正成功。这就是“至少一次”的本质没有收到消费者的确认时消息一定会被重新发送接收到确认后永远不会再发送同一批次的消息。如果你收到的确认发生在业务真正成功之前那么业务失败时消息也不会被重发。类比一下Kafka有点像快递柜它只管你取了件、签了字至于你签完字后包裹是不是真的用到正确的地方快递柜不关心。所以严谨地说Kafka能给你的只有位移提交这个钩子剩下的都需要你自己设计。2. 第一道防线把位移提交的时机彻底想清楚2.1 自动提交最大的坑不是“自动”很多人一听到自动提交就摇头其实自动提交本身不是魔鬼魔鬼是它的提交时机和你业务处理粒度对不上。enable.auto.committrue时Consumer 在每次poll()返回后下一次poll()调用前会自动提交上一次拉取到的所有消息的位移默认间隔是5秒。这个机制有两个隐患。第一如果一批消息还没处理完就发生了线程崩溃或poll()超时这批消息的位置还没提交重启后会从旧位置重新读一遍这是重复消费不算丢。第二如果业务处理过程特别耗时消耗的时间超过max.poll.interval.ms默认300秒Consumer 会被判定为失联从而触发rebalance把你踢出消费组而你刚好还没提交位移。被踢出去后你处理的结果已经写入了下游系统另一个消费者又从旧位移开始重放一遍重复写入就发生了。我见过一个典型案例消费者在处理一批消息时会批量把结果写入 Elasticsearch。这批操作需要在全部写完后一次性提交位移但中途某条数据导致写入报错代码没有对异常做正确处理也没有触发位移提交于是消费者重启后把整批消息都重新处理了一遍ES中出现了大量重复文档。问题的根源不是自动提交而是“提交时机不受控”。2.2 手动提交的正解先业务成功再提交位移关掉自动提交后最标准的姿势是poll一批消息逐条处理业务全部成功后调用一次 commitSync 提交位移。这样只要业务失败位移就不前进下次继续从失败位置开始读消息不会丢。但这个姿势也有个很隐蔽的坑如果下游系统不是幂等的而且业务处理成功、提交位移时抛了异常这条消息会被重复消费。对于不可幂等的操作比如发送短信、增加余额、扣减库存重复执行极有可能造成业务故障。所以严格说“先业务成功再提交位移”只是保证了消息能从Kafka侧读到并不能保证业务侧最终只执行一次。因此这里要引入一个概念业务结果和位移提交的一致性。你的提交时机必须基于“业务已经稳定落库”为前提任何可能影响业务成功的因素都必须先于提交处理好。如果下游是异步写的你得等到下游确认写成功而不是把消息丢进线程池就认为完成。2.3 一定要会用 commitSync 和 commitAsync 的组合commitSync 是同步阻塞式提交提交失败会抛出异常安全但慢。commitAsync 是非阻塞提交性能好但它是异步执行的调用后立刻返回如果提交失败没有自动重试机制只能通过回调感知一旦失败且进程继续往下消费这批位移就永久丢失消费者重启后会重新读取大量旧消息。我的线上标准玩法是“平时用异步关闭前用同步兜底”。具体操作是在消费循环里用 commitAsync 提交每次 poll 的位移保持消费性能在消费者关闭或者遇到退出条件时再调用一次 commitSync把最后一批位移同步落盘。如果这次 commitSync 失败要打印异常日志并中止关闭流程让运维人员介入检查而不是直接退出进程。很多生产事故都出在“只用 commitAsync 但从不检查回调结果”上不重视提交失败的后果就是消费者一重启几十万条消息重新砸向下游。3. 第二道防线消费端幂等把重复消费变为无害消费3.1 幂等不是Kafka帮你做的是业务自己的设计因为Kafka是“至少一次”所以重复消费无法避免但你可以设计业务逻辑让重复消费不产生错误影响这就是幂等。我在维护一个订单系统时消费者负责把MQ里的订单消息同步到搜索服务。初始版本没有做幂等某次升级时消费者重启导致过去两个小时的订单消息被重新推送了一遍搜索服务里的订单文档大量重复搜索排名和统计口径全乱了。当时做修复时我只加了一个唯一索引结果并发重放情况下依然出现了重复数据。原因是“先查再判断再插入”在并发时查到的都是旧数据判断结果都是“需要插入”于是又插入了一次。这个教训告诉我幂等不能靠业务代码里 if-else 判断必须靠数据库约束或条件更新来兜底。3.2 三种落地方式从简单到复杂第一种主键去重。消息体里带一个全局唯一业务ID比如订单号、流水号。下游表用这个ID建主键或唯一索引插入时捕获主键冲突。这种方式最简单适合只做“新增数据”的场景。第二种状态机去重。适合状态流转类业务比如支付单从 PENDING 到 PAID 再到 REFUND。你可以把更新SQL写成期望状态约束例如UPDATE payment SET statusPAID WHERE order_id? AND statusPENDING如果影响行数为0说明这条数据已经被处理过了直接用 return 结束。这种方式能避免重复消费覆盖掉更新的状态。第三种本地消息表加事务。适合业务需要同时操作多张表、或还要触发外部调用的场景。这也是我重点推荐的方式下一小节单独展开。3.3 本地消息表把消费和业务写入放进同一个事务本地消息表也叫 inbox 模式的核心思想是把消费行为本身变成一条数据库记录与真正的业务更新在同一个数据库事务里提交。具体做法是建一张消费记录表字段包含消息唯一ID、topic、分区、offset、消息体、处理状态。消费开始后先尝试往这张表插入一条记录同时在同一事务里执行业务更新。重复消费时因为消息ID已经存在插入主键冲突事务可直接回滚业务更新不会执行第二次。这个方案的关键是消费记录的插入和业务更新必须在同一个数据库事务里。如果业务更新成功了、消费记录没插上那业务已经写进去了重启后重放一遍又执行了一次更新幂等就失效了。反过来消费记录插上了、业务更新失败事务回滚两者都会回滚这个是正确的。我在订单系统里用这个方案后再也没出现过重启导致重复写入的问题。它的局限性也很明显要求业务数据库必须和消费记录表在同一个库里一旦你的系统是微服务架构、业务数据分散在不同库甚至不同中间件里这个方案就用不了。跨库跨系统的情况需要用到第5节的两阶段消费方案。4. 第三道防线别让处理太慢变成隐性消息丢失4.1 消费者被踢出组是很多线上事故的元凶Kafka消费组依赖心跳机制维护成员关系。如果消费者处理消息耗时太长超过max.poll.interval.ms默认的300秒Coordinator 会判定该消费者已经失联触发 rebalance把它的分区重新分配给其他消费者。这里有个致命细节消费者被踢出组后它并没有真的停止处理消息而是还在慢吞吞地处理自己之前 poll 到的那批消息。当它处理完后Kafka 要求它重新加入消费组重新 poll 时从已提交的位移继续读于是它之前处理但没提交的消息会被另一个消费者再处理一遍。如果另一个消费者已经处理完并提交了位移两边可能在同一个分区上互相覆盖位移导致整个消费组的进度来回横跳。我排查过一个真实案例某消费者在每次处理消息时都会调用第三方短信接口单条耗时超过2秒poll 了500条后处理完需要接近20分钟明显超过5分钟。于是这个消费者每隔20分钟就被踢出组一次日志里全是重平衡消息短信被重复发送了多次业务投诉大量到达。排查到最后跟Kafka本身没有关系完全是慢消费导致的。4.2 三个参数一起调才能解决慢消费问题解决慢消费单纯调大某一个参数是不够的我一般会把这组参数联动调整max.poll.interval.ms调大到能容纳最坏情况下一批消息处理时间的两倍以上。先估算单条消息的最大处理时间乘以max.poll.records得到最坏耗时再留出余量。max.poll.records调小默认500条对多数业务来说都太大建议按单条处理耗时调成50到100条这样整体处理时间可控。session.timeout.ms调整心跳超时时间但注意这个值不能设得太大否则故障检测变慢反而会掩盖消费者失联问题。通常配合heartbeat.interval.ms一起调心跳间隔要小于 session.timeout 的三分之一。举个例子假设单条消息最大处理耗时2秒一批调成100条最坏耗时200秒那么max.poll.interval.ms至少设成600秒才比较安全。如果业务里还有第三方调用还要额外考虑网络抖动、重试超时导致的耗时膨胀建议再乘以1.5的余量系数。4.3 终极解法poll线程和业务线程分离如果业务处理确实无法在几秒内完成或者下游经常抖动最稳妥的架构是让poll线程只负责拉取和提交位移业务线程池负责处理消息。poll线程拉到的消息先放入一个本地阻塞队列业务线程从这个队列取消息处理。这样poll线程可以保持心跳永远不会因为业务超时被踢出组业务线程慢点无所谓不影响组关系。这个方案有两个代价需要提前设计。第一个是本地队列会积压消息必须监控队列深度并设置上限超过上限就暂停poll防止内存溢出。第二个是进程崩溃时本地队列里的消息会丢失。所以这个架构必须配合“至少一次”的恢复机制poll线程提交位移的时机必须基于业务线程真正处理完成的确认不是poll完就提交。具体做法是业务线程处理成功后把对应的 offset 放入一个“已确认集合”poll线程在提交位移时只提交集合中已经确认的偏移量未确认的部分即使进程崩溃重启后也能从Kafka重新拉取。这就是自定义位点确认的雏形下一节展开讲。5. 真正的进阶玩法两阶段消费加外部位点做到不丢不重5.1 两阶段消费先把消息安全落盘再慢慢处理本地消息表方案要求业务和消费记录同库但在跨库、跨系统场景下这个前提不成立。我在生产环境用过一套更通用的模式这里叫它“两阶段消费”背后思想是把“拉取确认”和“业务处理”解耦成两个阶段第一阶段消费者 poll 到消息后先把消息完整写入本地的一张“待处理消息表”字段包括 topic、分区、offset、消息体、状态。消息安全落盘后消费者立即提交Kafka位移。此时从Kafka视角这批消息已经消费完成但你的本地表里还躺着这些记录状态是“待处理”。注意这里落盘操作不能失败否则位移已经提交、消息却丢了属于真丢。所以第一阶段最好用事务保护写文件或写数据库都要确保成功后再去提交位移。第二阶段后台任务一个独立线程池或定时任务从待处理消息表里按时间顺序捞记录真正执行业务处理逻辑。业务成功后更新记录状态为“已完成”或直接删除。这套模式的好处非常实用。poll线程永远不被业务阻塞不会触发rebalance业务可以随时重放即使第二阶段中途崩溃重启后捞出来的还是那些未完成记录跟Kafka位移毫无关系。我用了这套方案后排查消费链路健康状态也变得很简单直接查待处理表里的记录数和最老记录时间就能知道积压情况和进度比看Kafka Lag更直观。这个模式在行业里有时叫持久化收件箱模式但很多讲述Kafka的文章里很少把它作为消费端不丢消息的解法来介绍。它真正解决了“消费位移已经提交但业务还没执行完”之间的悬空状态。5.2 把消费位点同步到外部存储实现进度可观测、可回放Kafka自带的消费组位移存储在内部主题__consumer_offsets里。你没法方便地查看某个组每个分区消费到哪个 offset更没法手动把位移回退到几天前重放数据。于是我额外维护了一份外部位点结构很简单key 是consumer-group:topic:partitionvalue 是offset同时存上最近消费消息的时间戳存储用Redis或MySQL均可。同步方式是在手动提交Kafka位移时同时把当前批次的 offset 和时间戳写入外部存储。这个操作可以容忍延迟外部存储挂了不影响主流程下次启动时外部位点读不到就回退到Kafka自带位移继续消费。外部位点有两个实际价值。第一个是人工回放。某次上线后数据跑错产品要求把某个topic从3天前重新消费一遍如果没有外部位点你只能靠改groupId、改auto.offset.reset来“骗”出一条新消费链路非常麻烦。有了外部位点直接把Redis里的对应 key 回退到3天前的 offset重启消费者就能从那里开始读。第二个是跨环境恢复。测试环境想复现生产数据把生产外部位点导出在测试环境恢复消费者可以无缝从生产位置继续消费。需要提醒的是外部位点是用来扩展Kafka位移能力的不要用它替代Kafka位移。我的原则是Kafka位移始终作为兜底位点外部位点作为服务位点读取只有外部位点存在且校验通过时才覆盖Kafka位移否则一律用Kafka位移避免两边进度不一致把消费者搞蒙。5.3 Consumer的优雅停机不处理好会丢位移消费者被 kill 的时候位移怎么保这个细节很少有文章讲透。如果是kill -9强杀进程没话说shutdownHook 都不会执行。但如果你用kill -15优雅停机Java进程会执行 shutdownHook这时候你有机会把还没提交的位移补一次 commitSync。很多系统的坑在于 shutdownHook 里只做了关闭线程池、关闭 Consumer 这类动作但少了一步“最后提交一次位移”。Consumer.close()其实会提交位移但它提交的是最后一次 poll 之后、close 之前所有已消费的位移如果你的业务还没处理完这批消息close 不会帮你确认任何东西。最安全的做法是主动控制消费循环的退出条件退出循环后遍历“已确认集合”找到每个分区当前已确认的最大 offset手动提交一次 commitSync然后再关闭 Consumer。这样才能保证进程退出前所有已成功业务对应的位移都稳稳落盘。否则重启后就可能重放一大批消息如果没有幂等保护事故紧接着就来了。6. 消费端问题排查实录与常见误区6.1 一直在rebalance消费却一动不动症状消费者组日志里频繁打印 rebalance 信息消费进度不变业务侧也没有新增数据。排查思路先用命令行工具查看消费组状态如果状态是 Stable 但 Lag 不变多半是消费者和 coordinator 之间的会话有问题。最大嫌疑依然是max.poll.interval.ms触发也就是业务处理过慢。还有一种情况是客户端版本与服务端版本差异过大可能导致 coordinator 同步异常引发重平衡风暴此时先升级客户端版本到与服务端匹配的版本再按4.2的方法调参。6.2 位移提交抛 CommitFailedException消费却还在继续症状消费者日志连续报 CommitFailedException应用没有崩溃但已经能观察到重复消费数据。原因CommitFailedException 多数时候代表当前消费者已经不在消费组内了比如因为处理超时被踢出组。此时调 commitSync 是在无效提交因为它找不到当前消费者所属的分区。这种情况下继续消费实际是在“无主状态”下处理消息恢复后位移会被新的在位消费者覆盖数据就会重放。处理方式不要盲目重试提交。先判断异常类型如果是 RebalanceInProgressException等 rebalance 完成再提交如果是 CommitFailedException 且消息处理还在进行说明你的处理逻辑超时了应优先优化处理速度而不是在代码里死循环重试提交。6.3 下游偶发失败一重试就雪崩怎么办使用本地待处理表或两阶段消费模式时第二阶段处理失败怎么办我的经验是分三步走第一步单条重试N次比如重试3次每次间隔递增1秒、5秒、15秒 第二步超过重试次数后把这条记录标记为“失败待处理”不阻塞主链路让后台任务继续扫后面的记录 第三步定时任务每隔一段时间扫失败记录重新投递到死信队列或直接重放。核心原则是把坏消息隔离到一边不要阻塞整个消费链路。在主消费线程里用无限重试等待下游恢复是很多消费者被踢出组的直接原因。这条方法论不只是Kafka适用对RocketMQ、RabbitMQ同样成立本质就是失败隔离、定时补偿、幂等保护三件事。6.4 面试级回答消息不丢失的完整逻辑如果你准备面试或者要给团队做技术分享不建议机械地背“生产端acksall、消费端手动提交”。这套回答太单薄也没有体现体系。我建议按“生产端不丢、broker端不丢、消费端不丢”三层拆解消费端重点讲三件事第一位移提交时机必须和业务结果一致线上一律关闭自动提交 第二消费端默认至少一次语义必须配合幂等设计来吸收重复消费 第三慢消费者要用 poll 线程与业务线程分离、调整max.poll.interval.ms等方法防止被踢出组导致隐式重复。说完这三条再补一段自己的实践比如用本地待处理表加外部位点做重放就能明显比背书的人更有说服力。面试官想听的其实不是某个参数的默认值而是你有没有真正在生产环境里思考过消费端的可靠性和工程取舍。最后分享一个排查效率的小经验如果你在公司里没有现成的Kafka可视化平台全靠命令行查 lag 和消费组信息效率会很低特别是 topic 一多、消费者组一多的时候命令行输出的可读性很差。我习惯维护一个轻量级的 Kafka 监控后台把消费组的滞后量、分区分布、消费速率统一展示出来比整天敲命令直观得多。搭建时有一个坑监控工具默认会把所有消费者组的元数据都拉一遍集群 topic 数量大的时候页面打开非常慢。后来我把采集间隔调大只监控核心消费组和核心 topic才把页面响应速度提上来。像这种小工具不需要很复杂能帮你快速定位“哪个链路在积压、消费有没有卡住”就已经很大价值了。毕竟排查消息不丢的问题很多时候一半时间都花在找到问题出在哪个环节上。根据我个人的项目经验消费端保证消息不丢失本质不是靠某一个开关或参数而是靠一套组合设计正确的位移提交姿势、业务幂等、本地待处理表、外部位点回放、优雅停机兜底。这几个环节互相配合才能真正做到不丢不重。如果你已经在用本地消息表模式可以试着加一层外部位点如果只是简单改了个手动提交那至少把幂等补上。每次做完这类优化我都建议顺手写一份消费端架构说明把位移策略、幂等方案、异常处理流程记下来下次定位问题会节省大量时间。