
做Kafka的人迟早会遇到这么一幕消费组跑得好好的某天你重启其中一个实例它没有从头消费也没有跳过任何消息而是精准地从上次中断的位置继续。你可能会想这不就是Kafka的基本操作吗但真去追这个问题你会碰到两个绕不开的概念——位移offset和高水位high watermark。很多新手学Kafka最头大的也是这两个词位移还好理解高水位听着就像玄学网上资料东一榔头西一棒子越看越晕。这篇内容就是想把这两件事彻底讲透。我会从最底层的存储逻辑讲起解释位移到底存在哪、怎么提交、怎么越界再拆解高水位和LEO的关系、HW推进的完整链路、以及高水位机制在故障场景下引发的数据丢失问题。不管你是刚入门想搞懂原理还是已经在排查线上消费延迟和重复消费这篇都能给你一套能直接用的认知框架和实操思路。1. 位移的底层逻辑Kafka凭什么记得你上次读到哪1.1 分区位移与消费者位移别混为一谈先说个最常见的误区。很多初学者会把消息的位移和消费者的位移当成一回事。实际上Kafka里有两种位移它们虽然长得很像但完全不是一个东西。第一种是分区级的位移。每条消息写入某个分区后会获得一个从0开始递增的序号这就是这条消息在这个分区内的位移。比如分区里有100条消息它们的位移就是0到99。这个位移由生产者写入时分配一旦落盘就基本固定了代表的是消息在日志文件里的物理位置。第二种是消费组级的位移也叫消费位点。它记录的是当前这个消费组在这个分区上下一条要消费的消息位移是多少。举个例子一个消费组读完了分区位移0到49的消息那么它的消费位点就是50意味着下次从这里继续消费。这个区别为什么关键因为很多人把二者混在一起导致后续看Lag、调位移重置时完全懵掉。分区位移是数据本身的位置由Kafka存储层管理消费组位移是消费者读到哪了由消费端和协调器共同维护。Kafka能实现从上次位置继续读靠的正是消费组位移的持久化而不是重新扫描分区。1.2 位移存哪__consumer_offsets内部主题这个问题很多人第一次听到会愣一下Kafka的消费者位移居然也存在Kafka自己里没错。老版本的Kafka把位移存在ZooKeeper里后来因为ZooKeeper不适合高频写入、且和消费组协调逻辑耦合太重从0.9版本开始就改成了内部主题__consumer_offsets。这个内部主题也是用普通Topic的方式存的只是名字以双下划线开头Kafka默认不允许客户端直接对它生产消费。它默认有50个分区每个分区的副本数由broker端的offsets.topic.replication.factor配置决定如果broker数量不足会自动降为1。这里有个实操经验如果生产环境broker数量大于等于3一定要把这个副本因子显式设成3否则一旦broker宕机整个集群的消费组位移都可能丢失后果非常严重。那位移消息是怎么组织的每条位移提交记录它的Key是消费组ID 主题名 分区号的组合Value就是位移值和一些元数据。Kafka对这个Key做哈希算出它应该写到__consumer_offsets的哪个分区。所以同一个消费组、同一个分区的位移始终落在内部主题的同一个分区上。这也是为什么消费组协调器Group Coordinator可以高效地管理每个消费组的位移——它只需要处理内部主题的部分分区。1.3 初次消费从哪开始auto.offset.reset的细节与选择搞清楚了位移存哪另一个逃不掉的问题是一个全新的消费组或者位移已经过期/越界时Kafka该从哪里开始消费这由auto.offset.reset参数决定它有三个值earliest、latest、none。很多人以为earliest就是从第一条开始读latest就是从最新一条开始读这话大致没错但忽略了一个关键点这个参数只在没有已提交位移或提交的位移已失效时才生效。如果消费组已经正常提交过位移那么无论auto.offset.reset设成什么重启后都从提交的位移继续消费这个参数根本不会触发。earliest和latest的差异在实际运维中影响很大。比如你有一个离线分析任务消费组停了两天期间消息已经超过日志保留期限被删除了那么重启后位移就失效了。如果配置是latest它会直接跳到当前最新位置中间丢掉的旧数据不补如果配置是earliest它会尽量从最早的可用消息开始。所以对不允许丢数据的场景建议设earliest并搭配消息保留期评估对只关心当前增量的场景latest更合适。none则是严格模式如果找不到有效位移就直接报错适合用来暴露配置问题。2. 高水位和LEO一条消息什么时候才算对消费者可见2.1 高水位的水库类比可见性边界说完位移终于到高水位了。新手一听高水位三个字容易联想到水库或者水位线这个类比其实非常准确。你可以把Kafka的每个分区想象成一个水库生产者不断往里面注水水位就是分区里消息的写入情况消费者在岸边舀水它只能舀到水位线以下的水水位线以上的水虽然已经注进去了但它看不到。这里的水位线就是高水位High Watermark简称HW。更准确地说消费端能消费的最大位移是小于HW的也就是offset小于HW的消息才对消费者可见等于HW的消息还不能读。那为什么要设这么一条可见性边界直接让消费者读到最新写入的消息不行吗答案是不行。因为Kafka是一个分布式系统每个分区都有多个副本消息只有被足够多的副本同步之后才算真正安全。如果leader还没来得及把消息同步给follower就宕机了这条消息就会丢失。HW存在的意义就是把已写入但尚未充分复制的消息和已安全复制的消息区隔开防止消费者读到那些可能随时消失的数据。2.2 LEO怎么变化写入即增长但增长不等于可见与HW紧密相关的另一个概念叫LEO全称是Log End Offset日志末端位移。可以把它理解成当前日志写到哪了。生产者每成功写入一条消息对应分区的LEO就加一。所以LEO代表的是物理上已经存在的消息量而HW代表的是逻辑上允许消费的消息量。这里有个非常容易混淆的点生产者写入成功只代表LEO增加了不代表HW也增加了。HW的推进依赖于副本之间的同步进度通常它总是小于或等于LEO的。当HW等于LEO时说明所有满足条件的副本都已经追上了最新写入这时消费者才能读到真正的最后一条消息否则哪怕生产者已经确认写入消费者也会在HW处被挡住。我自己在排查消费延迟时见过不少这种情况消费组Lag为0但还是读不到最新消息查来查去发现是HW没有推进消费者被卡在高水位上。所以如果你只盯着消费组Lag而不看HW和LEO的关系很容易误判。2.3 HW的数学表达为什么是最小LEOHW具体等于多少取决于分区副本的同步情况。一个分区有一个leader副本和若干个follower副本。理想情况下如果所有副本都完全同步那么HW就等于所有副本的LEO也就是等于leader的LEO。但现实中follower的同步总有延迟于是HW的通用计算方式是取满足同步条件的所有副本LEO的最小值。这个满足同步条件的所有副本就是ISR即正在同步的副本集合。具体来说leader在推进HW时会查看ISR里所有副本的LEO然后取其中最小的那个作为新的HW。为什么要取最小值因为最小值代表着最慢的那个同步副本也已经到达的位置只有保证了这个位置HW以下的消息才不会因为某个副本掉队而丢失。这里就要注意了HW完全由最慢的那个副本决定。一个分区的ISR里有3个副本其中2个已经写完100条另外1个只写了30条那HW就是30。也就是说哪怕99%的消息都同步了只要有一个副本拖后腿消费者就只能读到前30条。这也是为什么集群运维中必须密切关注ISR膨胀和副本Lag——它们会直接压住HW进而表现为消费变慢。3. HW推进的执行链路从生产者写入到消费者可读之间发生了什么3.1 生产者写入之后leader的LEO先动要真正理解HW不能只背公式得把流程串起来。假设一个分区有3个副本一个leader两个follower。生产者把消息发到leader时leader会先把消息写入本地日志此时leader的LEO加一但HW保持不变。注意Kafka的写入并不是写leader就返回成功而是取决于acks参数。如果acks0生产者发出去就不管了消息可能写成功也可能没写成功没有任何确认。如果acks1leader写入本地日志后就会返回成功但此时follower可能还没同步。如果acksall或者叫acks-1生产者必须等待ISR中所有副本都写入成功后才返回。这种配置下生产者看到写入成功时消息其实已经安全复制了HW的推进会更快跟上。所以从消息生产到消费者可见中间是有时间差的。哪怕你用的是acksall也经历了两段延迟第一段是leader写入到follower同步完成第二段是follower同步完成到HW被推进。后者不太起眼但在副本Lag大的时候会直接拖慢所有下游消费。3.2 副本拉取与汇报机制follower的LEO如何影响HWfollower是怎么把数据同步过来的答案是主动拉取。follower会持续向leader发送FetchRequest把自己要拉取的下一条位移告诉leaderleader则从该位移开始返回一批消息。这里头藏着一个容易被忽视的机制follower在每次FetchRequest中都会带上自己当前的LEO。leader收到请求后不只是返回消息还会顺便根据这些信息更新它记录的、关于每个follower的LEO。然后leader会重新计算HW如果ISR中所有副本的LEO都超过当前HW就把HW推进到其中最小的LEO。而follower这边收到FetchResponse后也会做两件事第一把消息写入本地日志更新自己的LEO第二更新自己的HW。follower的HW更新规则是取响应里携带的leader HW和自身LEO的较小值。为什么不能直接用leader的HW因为follower可能复制得还不够快它的本地日志里根本没有leader HW以下那么多条消息如果直接采用leader的HW就会让消费者在某些支持follower读取的版本里读到不存在的消息。整个过程用大白话总结就是follower一边拉数据一边汇报自己走到哪了leader根据这些汇报决定水位线能放行多少。这个机制看起来简单但Kafka官方历次版本都在改它的细节原因就藏在后面的数据丢失问题里。3.3 为什么HW推进滞后会导致消费延迟聊完机制来看一个很实际的场景。假设你的Kafka集群有一个副本磁盘IO出现抖动同步速度骤降它的LEO远远落后于其他副本。因为HW取的是ISR中最小的LEO整个分区的HW就会卡在很低的位移上。此时上游生产一切正常生产者的acks只要不是all也不会感知到异常但下游消费者会明显感觉消息变慢了或者读不到最新的。这个现象有很强的迷惑性。很多运维同学一看到消费延迟上涨第一反应是加消费者并发结果发现一点用没有。其实瓶颈根本不在消费端而在HW推进上。判断方法也简单查看分区描述看ISR成员和副本Lag。如果ISR正常、Lag很小HW却迟迟不动就要去查broker端日志和JMX指标了。我自己处理过一个类似问题最后定位到是其中一个broker的Page Cache回收异常导致follower的fetch请求一直超时ISR里那个副本长期不在同步状态HW被压得死死的。重启那个broker后HW立刻恢复。这类问题在监控面板上往往表现为多个消费组同时Lag上涨如果监控没有分主题分broker的维度很容易大海捞针。3.4 ISR同步中副本的合格线上面多次提到ISR这里把它单独说一下。ISR全称In-Sync Replicas是正在同步中的副本集合。leader在推进HW时只考虑ISR里的副本不在ISR里的副本即使LEO落后再多也不会去拖低HW。那么怎么判断一个副本还在不在ISR关键参数是replica.lag.time.max.ms默认30秒。如果follower在这个时间内没有向leader发过FetchRequest或者发了但一直没有推进LEOleader就会把它从ISR中移除。反过来被踢出去的副本如果后续追上了leader末尾还能重新加回ISR。这里有个历史演变值得知道老版本里判断副本是否同步还有一个条件就是副本的滞后条数不能超过replica.lag.max.messages。这个配置后来被移除了因为它和同步的概念本身冲突——一个follower可以长时间不拉取但只要它在那30秒内拉了一次并追平就会被认为同步。改成时间维度后判断更符合实际。ISR变化对HW的影响我在生产环境里见过两个极端。一种是非ISR副本带病在ISR里HW被拖垮另一种是ISR频繁增删也就是所谓的ISR抖动导致HW忽高忽低消费端和监控端看到的Lag数据剧烈跳跃。排查ISR抖动通常要看broker端的网络、磁盘、GC指标以及follower的fetch线程是不是有异常的Full GC停顿。4. 位移提交的两种姿势自动提交的坑与手动提交的正确打开方式4.1 自动提交默认配置背后的定时器大多数人刚开始写Kafka消费者时用的都是默认配置其中enable.auto.committrue。它的工作机制很多人想当然地以为是每消费一条消息就提交一条其实不是。它是基于poll的消费者每poll一次后台的位移提交器会在间隔auto.commit.interval.ms默认5秒后把这一批poll返回的最大位移提交上去。自动提交最大的问题不是它不提交而是它提交的时机和业务处理时机是脱节的。举个例子你poll了一批消息还没来得及处理完5秒到了位移被提交成这批消息的最大位移。此时消费者进程突然宕机重启后从已提交的位置继续消费之前poll到但没处理完的消息就永久丢失了。反过来还有一种情况消息处理完了但还没到5秒的提交间隔进程挂了。重启后位移还是上次提交的旧值已经处理过的消息又要重新处理一遍。这就是自动提交的丢失和重复两个典型事故场景。所以如果你的业务对消息不允许丢也不能容忍大面积重复自动提交大概率不适合直接裸用。4.2 commitSync与commitAsync的取舍那手动提交该怎么做常见的是两个APIcommitSync和commitAsync。commitSync是同步提交它会阻塞当前线程直到Kafka确认提交成功才返回。好处是可靠提交失败能及时感知坏处是每次提交都带来一次RTT和一次刷盘等待吞吐会受影响。commitAsync是异步提交它发出提交请求后立即返回成功或失败通过回调通知。好处很明显不阻塞主流程消费吞吐能拉满。坏处同样明显异步提交不保证顺序可能上一次提交还没写完下一次提交就发出去了最终提交的位移可能是旧值覆盖新值造成消息重复。实操中我一般建议组合使用正常情况下用commitAsync保吞吐在消费者优雅关闭比如调用close之前时强制补一次commitSync把最后一次位移同步提交掉。这样既保住了大部分场景的吞吐又避免了进程退出时丢位移的尴尬。还有一点回调里不能只打日志要判断isSuccessful失败时要决定重试还是记录报警不能放任不管。4.3 从至少一次到精确一次的常用思路聊到位移提交就不可能绕开Kafka的消费语义。默认情况下Kafka提供的是至少一次At Least Once语义消息不会丢但可能重复。因为位移可能在处理后被提交一旦先提交后处理或处理后没提交都会造成重复或丢失。Kafka中的精确一次Exactly Once主要靠事务API实现代价是性能开销和实现复杂度不是所有场景都需要。如果你的业务场景允许幂等处理我强烈建议用幂等来解决问题而不是硬上事务。比如把消息的唯一ID写进数据库处理前先查重或者把处理结果写入外部存储时带上消息位移利用唯一键约束去重。这比引入Kafka事务要简单很多也更可控。如果一定要用事务API大致思路是这样把消费消息和提交位移放进同一个Kafka事务里同时把处理结果也写入Kafka。这样要么消息处理结果和位移一起提交成功要么一起回滚从根源上避免了处理成功但位移提交失败导致重复消费的经典问题。缺点是事务会显著增加生产端和消费端的延迟用之前先压测。4.4 位移越界OutOfRangeException与重置策略位移相关的问题里有一个报错出现频率很高OffsetOutOfRangeException。它表示消费者要去拉取的位移已经不在日志的有效范围内了。常见原因有两个一是消息过了日志保留时间被删除消费组位移还没来得及跟上二是消费组位移被误提交到了HW之外的位置比如某些消费框架的bug或手动seek导致。遇到这个异常先别急着改代码按顺序排查第一步查看当前消费组每个分区提交的位移和对应主题的log-start-offset、log-end-offset用kafka-consumer-groups.sh就能看到。第二步确认位移是比日志最老的消息还小还是比HW还大。第三步根据业务语义决定重置策略想重新消费就reset到earliest想跳过历史数据就reset到latest或者用--to-datetime精确重置到某个时间点。重置位移本身要非常小心。我见过不少事故是因为误操作reset了生产消费组导致大量数据重复处理。重置前要评估好当前Lag、消息保留期、下游处理的幂等性尽量在业务低峰操作并先在测试环境演练一遍。还有一点Kafka 2.x之后reset-offsets支持--execute参数不加这个参数只是干跑预览不会真正执行这个设计很友好建议每次先预览再执行。5. 高水位的暗面HW截断与Kafka的数据丢失往事5.1 副本故障恢复时的HW截断逻辑前面讲的都是HW作为可见性边界的光明面。接下来这部分是很多人不知道的暗面——高水位机制本身在极端场景下会带来数据丢失。理解这个问题才算真正懂高水位。先看一个基本场景一个分区有leader A和follower BA上的LEO和HW都是100B上的LEO和HW都是80。此时A宕机B被选为新的leader。B不会直接开始服务它会先把自己的日志截断到它的HW也就是80。为什么因为B不知道A在80到100之间有哪些消息为了保持数据一致性它只能选择遗忘那部分它没同步到的数据。这个截断操作本身是合理的但它暴露了一个问题A明明已经写入了20条消息80到100生产者也拿到了写入成功的确认现在这些消息却因为leader切换而彻底消失了。在acksall的配置下按理说只有ISR全部确认才会返回成功但如果A宕机时ISR里只有它自己和B而B又恰好滞后就会出现这种数据丢失。5.2 两个经典问题场景leader切换与HW截断丢失消息HW带来的数据丢失在Kafka 0.11之前有一个非常出名的复现路径很多老Kafka工程师都吃过这个亏。场景是这样分区有两个副本Aleader和Bfollower初始时A的HW是10B的LEO是10A写入了新消息LEO到12但由于B还没拉取A的HW保持在10此时B重启不是A宕机B会按照自己的HW10做截断然后向A请求拉取位移10开始的数据A返回消息并同步给B一切正常。但如果在B发送FetchRequest前A也宕机了B被选为新leader。此时B的HW是10它会截断自己的日志到10然后接受新的写入。等A恢复后A发现自己的HW是12而新leader B的HW是10于是A也会被要求截断到10。结果就是那两条位移10和11的消息在所有副本上都被删掉了。这个问题的根源在于A在推进HW之前只参考了自己的LEO和滞后副本的LEO但A记录的HW10并不能准确反映哪些消息已经在所有副本上安全同步。因为A知道自己还有数据没同步出去但B不知道A还有高于HW的消息。旧版Kafka用HW作为截断依据的机制缺乏一个全局的协调信息导致双方的高水位互相不信任。5.3 leader epoch给位移加上版本号Kafka是从0.11版本开始引入leader epoch机制来修复HW截断这个致命问题的。leader epoch可以理解成leader任期或者说给每次leader变更加上一个递增的编号。每个broker都会持久化一张映射表记录第几任leader任期从哪个位移开始接管。有了这个起始位移信息之后新leader处理旧leader的follower恢复时就不是简单粗暴地按HW截断了。新leader会告诉follower我在第N任任期是从某个位移开始的如果你在上一任任期的数据领先于这个起始位移那些多出来的数据其实是不可信的应该截断。反过来如果follower的位移落后于新leader的起始位移就正常从那个位置继续拉取。这套机制本质上是把数据是否可信的判断依据从高水位是多少升级成你这批数据是哪任leader写的、合法性如何。HW仍然负责定义消费者可见性但不再充当日志截断的唯一依据。这算是Kafka历史上一次里程碑式的修复。如果你还在用0.11之前的版本看到这里应该明白升级到新版本不只是加功能更是在消除这种隐性丢数据风险。6. 实战排查位移不推进、HW卡住时我做了什么6.1 三分钟定位消费组Lag命令行工具箱理论讲再多最后都得落到排查能力上。先搭一套最基础的命令行工具箱我平时定位消费组问题80%都能靠这几个命令完成。查看消费组当前位移和Lag用kafka-consumer-groups.shkafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --group my-consumer-group --describe它会输出每个分区的当前位移CURRENT-OFFSET、日志末端位移LOG-END-OFFSET以及两者的差值Lag。看到Log-End-Offset长时间不动而生产端明明有流量说明HW或LEO出了问题。如果想看主题级别的消息范围可以用GetOffsetShellkafka-run-class.sh kafka.tools.GetOffsetShell \ --broker-list localhost:9092 --topic my-topic --time -1--time -1表示取最新位移--time -2表示取最早位移。结合这两个值就能判断日志最早的消息位移是否已经大于消费组位移确认是不是发生了位移越界。排查位移问题时我习惯把这两条命令的输出并排放在一起比对一眼就能看出异常在哪。6.2 HW相关的JMX指标与监控命令行适合即时排查但线上问题最好靠监控提前发现。Kafka broker会暴露大量JMX指标和高水位直接相关的是每个主题分区Log下的HighWatermark属性。在JMX里路径大致是kafka.log:typeLog,nameHighWatermark,topicxxx,partition0。同时可以看kafka.server:typeFetcherStats下面的指标了解followerfetch线程的请求成功率。真正该监控的不是HW的绝对值而是它的变化率和与LEO的差值。如果某个分区的HW长期小于LEO说明副本同步进度不正常。这时候就要配合查看ISR变化指标比如kafka.server:typeReplicaManager,nameIsrShrinksPerSec和IsrExpandsPerSec这两个指标能告诉你ISR在频繁增减是排查副本抖动的重要线索。对这些JMX指标我建议至少配置两个维度的告警一是ISR收缩或扩增速率的突增二是某个分区HW与LEO差值持续超过阈值。后者我一般以分钟级持续超过1分钟为告警条件避免短时间抖动误报。同时告警一定要带上broker和topic维度否则一个大集群里某个分区出问题聚合后的指标看起来毫无波澜根本看不出来。6.3 几个真实案例ISR抖动、磁盘故障、fetch线程阻塞最后分享几个我实际处理过的案例都是高水位和位移问题在真实环境里的表现希望能帮你建立直觉。第一个是ISR抖动导致的消费延迟。现象是某几个消费组的Lag每隔几分钟就跳一次监控图上像锯齿。排查后发现其中一个broker的磁盘延迟时好时坏follower的fetch请求经常超过replica.lag.time.max.ms被踢出ISR过一会儿又追上拉回ISR。每次踢出再拉回HW都会重新计算消费延迟就跟着波动。最终是更换掉异常的磁盘盘位解决的期间临时调大了replica.lag.time.max.ms把抖动的影响压住。第二个是磁盘读写故障引发的HW停滞。某个分区的follower一直无法写入本地日志LEO卡住不动但它在30秒内仍然持续发送fetch请求所以没有被踢出ISR。这就导致leader侧认为它还在同步HW被它的LEO死死拖住。这个场景的隐蔽性在于ISR看起来正常Lag也不是很大但消费端就是觉得消息延迟。最后是通过对比所有副本的Log End Offset才定位到问题那个故障副本的LEO明显落后。所以排查HW不推进时一定要逐副本看LEO不能只看ISR。第三个案例是消费者端fetch线程被Full GC打断。这是一个Java消费者进程大量对象分配触发频繁Full GC导致poll间隔拉长位移提交和消息拉取都变慢消费组频繁触发rebalance。rebalance本身不会重置位移但会让消费组的broker侧状态反复调整监控上看起来就像位移和Lag都在乱跳。处理方式是把消费端的堆内存调大、减少不必要的对象分配并调整了GC参数问题才消停。这个案例提醒我位移问题不一定是Kafka服务端的问题客户端自身的运行状态同样会表现得非常诡异。如果你正好处在排查现场我的建议是先冷静梳理位移提交、HW推进、副本同步这三条线分别确认状态再交叉定位。Kafka的很多问题表面上是消费延迟根子却在副本同步或者客户端GC上一步到位正好命中问题的情况其实不多。做Kafka排查做了这么些年我个人的体会是位移和高水位这两个概念属于看着简单、用着容易翻车的类型。位移的坑主要在消费端高水位的坑主要在服务端副本管理而两者一旦相互作用比如副本故障导致HW后退进而让消费组位移越界排查难度会成倍增加。这篇内容里的场景和分析都是我实际踩过或者帮别人排查过的希望能让你少走点弯路。以后不管是看面试题还是处理线上故障再遇到HW和LEO的时候至少心里有一条清晰的调用链路知道该往哪个方向查了。