为了消除 Lag团队把订单 Topic 从 12 扩到 36 个分区。Kafka 吞吐上去了订单状态却开始回退数据库锁等待也翻倍。扩分区同时改变了 Key 路由和最大并行度而旧数据不会自动搬家。扩分区不是无损扩容它不可缩回可能让同 Key 新旧记录跨分区还会把下游并发上限一起放大。三个同时发生的变化1. Key 重映射默认 Key 路由依赖序列化后的 Key 与当前分区数。分区数改变后同一 Key 的新记录可能进入不同分区历史记录仍留在旧分区跨分区没有总序。2. 历史数据不重分布新分区从创建时开始接收数据。旧分区上的热数据、磁盘占用和 Lag 不会被自动摊平因此“扩三倍”不等于现有热点立即下降三倍。3. 并发上限提高Consumer Group 可以同时激活更多分区任务。若每个任务都持有数据库连接或调用同一接口下游会在再均衡后突然承受更高并发。Kafka 官方运维文档把增加分区作为变更操作并明确警告不要手工增加内部状态 Topic 的分区分区数不能用同一命令缩减。Basic Operations为什么“先扩了再说”没有可靠回滚方案是否恢复旧路由是否保留新写入代价把分区数改回去不可行—Kafka 不支持缩分区停掉新分区 Consumer否新分区积压只止住下游压力自定义旧映射仅对新消息可能需处理新分区历史易形成双轨新建 Topic 重分区可以设计需迁移与切换成本最高但可控所以恢复点不能是“原分区数”而应是扩容前准备好的新 Topic/双写/读取切换方案。变更前只读评估bin/kafka-topics.sh --bootstrap-server broker:9092\--describe--topicorder-events bin/kafka-consumer-groups.sh --bootstrap-server broker:9092\--describe--grouporder-service第一条确认分区、副本和 ISR第二条确认 Lag 是否均匀。若 Lag 只集中在一个热分区增加空分区通常无效。还要离线抽样真实序列化 Key用旧、新分区数计算映射变化比例检查业务是否依赖 Key 内顺序、自定义 Partitioner 是否读取分区数、下游允许的最大并发。更安全的选择顺序先优化单分区处理慢调用、批次、压缩、下游写入。若是热 Key评估能否引入可聚合的子 Key不能破坏顺序时接受该 Key 的单分区上限。若只是未来容量不足且无顺序约束可直接扩分区但先限制 Consumer 并发。若既要扩容又要保持 Key 路由创建新 Topic固定新分区策略双写并校验后切读。高风险执行边界实际增加分区的命令会改变集群状态bin/kafka-topics.sh --bootstrap-server broker:9092\--alter--topicorder-events--partitions36不要把它当排查命令。执行前必须审批精确 Topic保存 Topic 配置与 Key 映射样本确认不是内部 Topic设置 Producer/Consumer canary限制下游连接和并发定义成功、停止和迁移条件。成功条件应同时包含分区级生产/消费吞吐提升、顺序违规为 0、下游错误与锁等待未恶化。出现业务版本回退、热点未改善或下游饱和即停止放量已创建的分区保留按预案切换新 Topic 或限制其使用。扩容后的验证按 Key 检查业务版本是否单调定位是否跨旧、新分区。比较每分区消息率与 Lag确认新增分区真的承接流量。核对 Consumer 活跃任务数和数据库连接/锁等待。对旧分区持续观察直到历史积压消化不能只看总 Lag。技术验收是路由、吞吐和错误率业务验收是订单状态不回退、同一订单动作不重复且端到端延迟达标。源码与 Java变更前先计算 Key 映射再等待 Admin Future以下源码定位与 Java 示例按 Kafka 4.3.1 静态审阅未在本环境运行示例包含不可逆的分区扩容操作只能用于经过审批的隔离 Topic。Key 路由由BuiltInPartitioner完成Admin.createPartitions进入KafkaAdminClient的请求链。以下代码会改变 Topic必须只在隔离环境执行。importjava.util.*;importorg.apache.kafka.clients.admin.*;publicclassPartitionExpansion{publicstaticvoidmain(String[]args)throwsException{PropertiespnewProperties();p.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG,localhost:9092);try(AdminadminAdmin.create(p)){Stringtopicorder-events-test;intoldCountadmin.describeTopics(List.of(topic)).allTopicNames().get().get(topic).partitions().size();intnewCountoldCount*2;System.out.printf(MUTATING topic%s partitions%d-%d%n,topic,oldCount,newCount);admin.createPartitions(Map.of(topic,NewPartitions.increaseTo(newCount))).all().get();System.out.println(changedtrue; rollback-by-shrinkfalse);}}}映射是increaseTo → KafkaAdminClient → Controller 元数据变更 → Producer 刷新 metadata → BuiltInPartitioner 新映射。先离线抽样 Key再运行变更代码不能回退分区数也不证明下游能承受新增并发。结论分区既是 Kafka 并行单位也是顺序边界和下游并发放大器。扩容前必须证明瓶颈确在分区数并把 Key 重映射、历史不迁移、不可缩减和下游容量写进同一份变更方案。