为了消除 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{Propertiesp=newProperties();p.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG,"localhost:9092");try(Adminadmin=Admin.create(p)){Stringtopic="order-events-test";intoldCount=admin.describeTopics(List.of(topic)).allTopicNames().get().get(topic).partitions().size();intnewCount=oldCount*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("changed=true; rollback-by-shrink=false");}}}映射是increaseTo → KafkaAdminClient → Controller 元数据变更 → Producer 刷新 metadata → BuiltInPartitioner 新映射。先离线抽样 Key,再运行变更;代码不能回退分区数,也不证明下游能承受新增并发。
结论
分区既是 Kafka 并行单位,也是顺序边界和下游并发放大器。扩容前必须证明瓶颈确在分区数,并把 Key 重映射、历史不迁移、不可缩减和下游容量写进同一份变更方案。