22-副本机制与ISR详解

发布时间:2026/9/6 4:32:27
22-副本机制与ISR详解 第22讲 | 副本机制与 ISR 详解导读Kafka 的高可用不靠多数派写而靠 ISR——一份动态收缩的同步副本集合。理解 Leader/Follower、LEO/HW 与滞后即踢出的机制你才能真正回答acksall 到底等多久、等谁。本讲目标理解副本的角色划分Leader / Follower、OAR / ISR / OSR彻底搞懂 LEO 与 HW 两个关键水位的定义、更新时机与用途掌握 Follower 同步的拉模式全流程与截断truncate逻辑理解副本被踢出/加入 ISR 的判定条件与相关参数能用命令与指标排查ISR 收缩消息看不到HW 不动等生产问题。一、为什么需要副本单点与持久性的矛盾分区是 Kafka 并行与容错的最小单位。一个分区只在一个 broker 上是单点副本replica把同一份数据放到多个 broker 上。KRaft 模式下副本放置由 controller 决定尽量机架/节点分散。Leader 副本唯一对外服务的副本处理该分区所有读写客户端元数据请求里返回的就是 Leader 所在节点Follower 副本不服务客户端读写唯一的职责是向 Leader 拉取消息保持同步随时准备在 Leader 挂掉时顶上。注意两点反直觉的设计消费者只读 LeaderFollower 不提供读扩展3.x 之前的follower fetching只是消费者侧的可选优化broker 端副本仍不服务生产者写。这与很多系统读写分离的思路不同原因在第20讲已铺垫——读写都在 Leader 上PageCache 局部性与零拷贝才能生效而一致性开销读 Follower 要处理读到比自己旧的 HW的问题不值得Follower 是拉pull不是推pushFollower 主动向 Leader 发FETCH请求。这复用了消费协议的代码路径也让 Follower 自己控制追赶节奏。几个集合名词一次讲清名词含义OAROrdered Assignment of Replicas分区的全部副本列表含顺序如[0,1,2]表示 broker 0 是首选ARAssigned Replicas同 OAR常见中文写法分配副本集合ISRIn-Sync Replicas与 Leader 保持同步的副本子集动态变化永远包含 Leader 自己OSROut-of-Sync ReplicasAR 中被踢出 ISR 的滞后副本AR ISR ∪ OSR恒等式二、LEO 与 HW一切同步语义的两个水位2.1 定义LEOLog End Offset日志末端位移即下一条待写入消息的 offset。LEO8 表示已有 offset 0~7 共 8 条消息HWHigh Watermark高水位所有 ISR 副本都已复制的最大位移 1即消费者可见的消息上界不含。消息只有进入 HW 之下才对消费者可见。为什么要 HW因为消费者不能读到只有 Leader 有、Follower 还没有的消息——否则 Leader 一挂、新 Leader 没有,这些消息就从已消费变成不存在破坏一致性。2.2 一个消息爬上 HW的完整过程设分区有 3 副本Leader ALEO5Follower B、CLEO 各 4当前 HW4。渲染错误:Mermaid 渲染失败: Parse error on line 15: ...5 Note over A,B,C: HW5offset4 对消 ----------------------^ Expecting TXT, got ,要点拆解Follower 的 FETCH 请求带上fetchOffsetLeader 由此得知该 Follower 的 LEO并更新其远程副本 LEOremote LEOLeader 每次处理 FETCH/PRODUCE 后重新计算HW min(Leader自身LEO, 各ISR Follower的remote LEO)HW 的传播本身也是滞后的Follower 要到下一轮 FETCH 的响应里才能拿到新 HW因此存在消息已同步但消费者还看不到的短暂窗口acksall的语义是写入被 ISR 中全部副本收到即成功严格说是 Leader 等待消息进入 ISR 全部副本的 LEO后应答它不等于消息对消费者可见——可见还要等 HW 推进。2.3 Leader Epoch修补 HW 机制的窟窿纯靠 HW 做截断有一致性风险。经典事故场景Follower B 落后HW2Leader A 的 LEO4A 宕机B 被选为新 LeaderB 把自己截断到 HW2继续接受新消息offset 2、3 被覆盖A 重启回来做 Follower按老机制截断到 HW——但如果它记得自己曾有 offset 3而新 Leader 的 3 是另一条数据就出现数据分叉或丢失。Kafka 0.11 引入Leader Epoch领导纪元解决每次 Leader 变更epoch1每个段目录维护leader-epoch-checkpoint文件记录每个 epoch 对应的 LEO 起点。Follower 重启/追赶时不再问HW 是多少而是先发OffSetsForLeaderEpoch请求问 Leader我在 epochX 时你有 offset Y你现在这个位置对应的 epoch 和 LEO 是什么然后精确截断到纪元切换点避免基于陈旧 HW 的错误截断。3.x 中 epoch 机制已是默认且不可关闭。三、Follower 同步与踢出 ISR的判定3.1 同步的判定参数replica.lag.time.max.ms 30000 默认 30 秒判定逻辑Kafka 0.9 改为基于时间而非条数Leader 为每个 Follower 维护它最后一次 caught-up 的时间戳即该 Follower 的 fetchOffset 追上 Leader LEO 的时刻若某 Follower 在replica.lag.time.max.ms内没有追上过Leader 的日志末端Leader 就把它移出 ISR反过来被踢出的副本只要重新追上fetchOffset Leader LEO就会被重新加入 ISR。为什么从落后条数改为落后时间旧机制replica.lag.max.messages在流量突增时会误杀正常副本瞬间涌入十万条大家全都超条数流量低谷时又检测不出真故障。时间维度对流量不敏感更能表达这个副本还活着且跟得上。3.2 ISR 收缩的代价可靠性与可用性的跷跷板ISR 收缩到只剩 Leader 时acksall退化为acks1的持久性min.insync.replicas 不设防的话min.insync.replicas2acksall是至少两副本落盘才允许写入的标准生产组合ISR 数 2 时生产者收到NotEnoughReplicasException牺牲可用性保住不丢ISR 变动会触发元数据变更ZK 模式写 /controller/isr_changesKRaft 模式进元数据日志并推送给客户端所以高频抖动的 ISR 也会带来集群层面的开销。3.3 Unclean Leader Election最后的开关当 ISR 只剩 Leader 且 Leader 宙机分区不可用。unclean.leader.election.enabletrue允许从OSR落后副本中选新 Leader代价是落后部分的消息永久丢失旧 Leader 恢复后要截断对齐。金融类业务必须保持false默认可观测性日志类可酌情开启。四、用一张图串起副本生命周期副本分配/新Leader当选fetchOffset 追上 Leader LEO超过 replica.lag.time.max.ms 未追上重新追上 LEO被选为 LeaderController 决策宕机后恢复按 Leader Epoch 截断再同步分区重分配/缩容FollowerSyncISROSRLeaderFollowerDeleted实战案例观测 ISR 收缩与 HW 推进命令演示# 1. 创建 3 副本 topicbin/kafka-topics.sh --bootstrap-server localhost:9092\--create--topicpay-events--partitions3--replication-factor3\--configmin.insync.replicas2# 2. 查看副本分布与 ISRIsr 列即当前同步副本所在 brokerbin/kafka-topics.sh --bootstrap-server localhost:9092\--describe--topicpay-events# Topic: pay-events Partition: 0 Leader: 1 Replicas: 1,2,0 Isr: 1,2,0# 3. 模拟故障停掉 broker 2持有副本的节点观察 ISR 收缩# 观察 Isr 逐渐变为 1,0且生产者仍可用因为 min.insync.replicas2bin/kafka-topics.sh --bootstrap-server localhost:9092--describe--topicpay-events|head-1# 4. 重启 broker 2观察它重新进入 ISR# 日志关键字Shrinking ISR / Expanding ISRtail-f/path/to/kafka/logs/server.log|grep-EShrinking|Expanding小型 Java 演示观测写入成功与消费者可见的 HW 差importorg.apache.kafka.clients.consumer.*;importorg.apache.kafka.clients.producer.*;importorg.apache.kafka.common.TopicPartition;importorg.apache.kafka.common.serialization.StringDeserializer;importorg.apache.kafka.common.serialization.StringSerializer;importjava.time.Duration;importjava.util.*;/** * 第22讲实战观测 endOffsetsLEO 视角与 HW 推进。 * 场景秒杀系统写单后立刻查询已下单数偶发少读——本程序复现并解释 * acksall 写成功与消费者立刻能读到之间存在 HW 推进的窗口。 */publicclassHwLagDemo{privatestaticfinalStringBOOTSTRAPlocalhost:9092;privatestaticfinalStringTOPICseckill-orders;publicstaticvoidmain(String[]args)throwsException{// ---------- 生产者acksall单条发送并等待确认 ----------PropertiespnewProperties();p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,BOOTSTRAP);p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,StringSerializer.class.getName());p.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,StringSerializer.class.getName());p.put(ProducerConfig.ACKS_CONFIG,all);try(KafkaProducerString,StringproducernewKafkaProducer(p)){RecordMetadatamdproducer.send(newProducerRecord(TOPIC,sku-A,{\orderId\:\S001\})).get();System.out.printf(写入确认: partition%d offset%d acksall 已返回%n,md.partition(),md.offset());}// ---------- 消费者立刻查询并轮询展示可见性延迟窗口 ----------PropertiescnewProperties();c.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,BOOTSTRAP);c.put(ConsumerConfig.GROUP_ID_CONFIG,seckill-counter);c.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class.getName());c.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class.getName());c.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,latest);try(KafkaConsumerString,StringconsumernewKafkaConsumer(c)){TopicPartitiontpnewTopicPartition(TOPIC,0);consumer.assign(List.of(tp));consumer.seekToEnd(tp);// endOffsets 返回的是下一条待写 offset即该分区 LEO 视角longleoconsumer.endOffsets(List.of(tp)).get(tp).offset();System.out.printf(分区 LEO%dendOffsets 视角%n,leo);// 轮询最多 2 秒观察消息何时越过 HW 变为可见longdeadlineSystem.currentTimeMillis()2000;booleanseenfalse;while(System.currentTimeMillis()deadline!seen){for(ConsumerRecordString,Stringr:consumer.poll(Duration.ofMillis(200))){System.out.printf(消费者读到: offset%d value%s%n,r.offset(),r.value());seentrue;}}System.out.println(seen?消息已越过 HW可见:2 秒内未读到HW 尚未推进);}}}关键点解读send().get()返回只代表ISR 全部收到不等于消费者可见——HW 的推进依赖 Follower 的下一轮 FETCH通常这个窗口只有几毫秒Follower 拉取很勤但在 Follower 忙/网络抖动时可能放大这就是写后立刻读偶发少读的原理性解释若业务要求写后立刻读自己写的数据正确姿势是用返回的partition/offset做点查校验或者读写分离到别的存储而不是加大轮询。踩坑提示 / 生产建议标准生产组合replication.factor3min.insync.replicas2 生产者acksallenable.idempotencetrue。低于这个组合请自觉评估丢失场景。min.insync.replicas 不要等于 replication.factor任何一个副本抖动都会让分区不可写可用性极差也不要设 1失去保护意义。replica.lag.time.max.ms 调优跨机房/高负载集群可适当调大到 60s减少 ISR 抖动调得越大故障切换时丢数据窗口仅 unclean 场景与元数据滞后越大。监控三个指标UnderReplicatedPartitionsISR 收缩数0 即告警、IsrShrinksPerSec、IsrExpandsPerSec。频繁 shrink/expand 说明磁盘或网络有问题别只盯着 CPU。消费者 lag 与副本 lag 是两回事前者是消费组落后后者是 Follower 落后监控面板别混在一页误判。扩缩容/重平衡期间 ISR 抖动是正常的别在 rebalance 窗口期给 UnderReplicated 告警打电话。本讲小结副本机制是 Kafka 可用性的基石Leader 独扛读写Follower 以 FETCH 拉模式追赶ISR 是时间窗口内追得上的动态白名单落后即踢、追上即回HW 由 ISR 内最小 LEO 决定划出消费者可见边界Leader Epoch 则修补了基于 HW 截断的一致性漏洞。写路径的可靠性语义acksall min.insync.replicas完全构建在这套机制之上。思考题ISR 已经缩小到只剩 Leader 一个副本此时acksall的写入还会成功吗它和acks1有区别吗结合min.insync.replicas讨论两种配置下的行为差异。为什么 Kafka 选择消费者只读 Leader而不是读最近的 Follower如果读 Follower需要额外解决什么问题提示HW 语义、会话黏性。