Kafka 运行链:Partition、Replica、Offset 与 Rebalance 如何协同
从追加日志、分区副本、生产批次、消费位点和组协调器出发,解释吞吐、局部顺序与故障恢复的真实代价。
Kafka 官方 Protocol Design 将 partition 定义为有序提交日志,并说明 key 驱动的语义分区;Consumer Configs 定义消费组协议与心跳边界。
Partition 同时承担并行度、顺序域和故障域
Kafka Topic 被拆成多个 partition,每个 partition 是只有一个 leader 接受追加的有序日志。生产者可显式指定 partition,也可由 key 的序列化字节经过分区器决定。相同 key 在分区数量与算法不变时进入同一分区,从而获得局部顺序;无 key 的记录更偏向批量和负载均衡。所谓 Kafka 全局顺序,除非 Topic 只有一个 partition,否则通常不存在。
增加 partition 会改变 key 到 partition 的映射,并且历史记录不会重排。依赖按用户、订单或聚合根顺序处理的系统,需要把扩容视为数据模型迁移:接受新旧记录跨分区、引入逻辑序号,或切换到新 Topic。把 partition 数只当吞吐旋钮,会在扩容时破坏状态机。
acks、ISR 与幂等生产者解决的是日志写入问题
acks=all 要求 leader 等待当前同步副本集合满足写入条件;min.insync.replicas 与副本数共同决定在故障时继续写还是拒绝写。acks=1 允许 leader 本地追加后响应,若其在复制前失效可能丢失已确认记录。生产者批次、linger、compression 和 batch.size 决定吞吐与延迟,消息过大则会同时挤压网络、页缓存和副本追赶。
幂等生产者用 producer id、epoch 和分区内序号抑制客户端协议重试产生的重复;它不能识别应用层重新构造并发送的同一业务事件。事务生产者可原子写多个 Kafka partition,并把消费位点纳入同一 Kafka 事务,但不能原子提交外部数据库。消费者还需 read_committed 才会屏蔽已中止事务。因此“Kafka exactly-once”是特定 Kafka 读写拓扑的能力,不是跨任意存储的宇宙事务。
Offset 是下一条读取位置,不是业务完成证明
消费组为每个 TopicPartition 保存已提交 offset。处理完 offset 9 后应提交 10,表示下次从 10 开始。先提交 offset 再写数据库会静默丢业务;先写数据库再提交 offset 会在崩溃时重复处理。后者配合幂等键、唯一约束或状态版本更可恢复。
Rebalance 会撤销旧成员的 partition 所有权并分配给新成员。旧消费者如果在撤销后仍提交或写入,会覆盖新消费者进度或产生并发更新。处理循环要限制 max.poll.interval 内的工作,长任务应拆分、暂停拉取或转交任务系统;撤销回调要停止新工作、等待有限在途任务并提交可证明完成的位点。
用所有权转移图读懂可靠性
失败矩阵比“保证不丢”更可执行
工程选择的核心是把不可避免的不确定窗口转成可重复、可查询、可补偿的状态。不要用无限重试掩盖未知结果,也不要用提前确认换取表面低 lag。
副本追赶与保留共同决定可恢复窗口
副本落后是否仍属于 ISR,取决于复制进度与时间阈值。leader 切换只从合格副本中选举时,系统宁可暂时不可写也避免选择落后副本;允许非同步副本选主则用数据一致性换可用性。生产者的 acks 不能脱离副本数、min.insync.replicas、非同步选主和机架分布单独评价。
保留策略按时间或大小删除日志段,compaction 则按 key 保留较新值并处理 tombstone。消费组停机超过保留窗口后,即使 offset 仍在,也可能无法读取所指数据;恢复只能从现存位置或外部快照重新构建。状态型消费者应把快照、changelog、重放速度和 schema 兼容一起设计,不能把“Kafka 可回放”理解为永久保留。
用两个 Java 状态模型固定不变量
下面的模型不连接真实 Broker。它们把本篇最容易混淆的状态、位置或所有权变化压缩成确定输出;随后再用 RabbitMQ、Kafka 或 RocketMQ 的集成环境验证协议、持久化、重平衡和故障时序。
javac --release 17 -Xlint:all -Werror examples/backend-development/message-event/kafka-runtime/PartitionOffsetDemo.java examples/backend-development/message-event/kafka-runtime/RebalanceOwnershipDemo.java
java -cp examples/backend-development/message-event/kafka-runtime PartitionOffsetDemo
java -cp examples/backend-development/message-event/kafka-runtime RebalanceOwnershipDemopartition=1 processedOffsets=[7, 8, 9] committedNextOffset=10 replayFrom=10
epoch=4 revoked=[orders-1] staleOwnerCommitAccepted=false assignedTo=consumer-B模型输出变化时,应定位是哪条业务不变量被修改,而不是只更新示例字符串。真实集成测试还要注入进程崩溃、响应丢失、Broker 切换和下游超时,因为这些窗口无法由纯内存模型模拟。
契约、权限与数据生命周期不能交给默认值
用分层测试证明协议与业务同时成立
把设计决定写成可以被证伪的承诺
指标、日志与容量要围绕状态迁移
生产侧观察发送尝试、确认延迟、未确认在途、返回/未路由、超时阶段和批次字节;Broker 侧观察 ready、unacked、partition lag、最老年龄、磁盘与副本健康;消费侧观察拉取、处理成功、业务拒绝、可重试失败、死信、重复命中、提交延迟和在途任务。吞吐必须与 p95/p99 时延、错误率、积压年龄一起看。
容量模型至少计算平均/峰值每秒消息数、平均/p99 字节、保留时间、副本倍数、重试放大、批处理密度、消费者处理时长与下游并发。恢复测试要在稳定生产流量仍存在时注入积压,验证净排空速率,而不是暂停生产后测一个理想峰值。
上线与演练以可恢复为验收标准
故障演练依次注入生产响应丢失、Broker 节点切换、消费者提交后崩溃、下游超时、永久坏消息和热点 key。每次演练都回答三件事:业务不变量是否保持、重复或延迟是否有持久证据、系统是否在容量预算内自动或人工收敛。
