这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 主题与分区是 Kafka 的“表结构”,一旦上线就很难随意修改。本章讲清楚分区数怎么定、副本因子怎么配、以及用 AdminClient 做自动化管理。
7.1 分区数怎么规划
分区数影响三个维度:
- 并行度:消费者实例上限 = 分区数;
- 吞吐:分区分布在多个 Broker 上,写入与读取并行;
- 资源与开销:每个分区都是一个目录、多组文件、多个副本、多个网络请求对象。分区不是越多越好。
经验公式
目标吞吐 / 单分区吞吐 = 所需分区数
例:目标 100 MB/s,压测单分区 10 MB/s -> 至少 10 个分区
考虑峰值与增长,规划为 20-30 个
再检查消费端:如果下游每秒只能处理 1 万条、单消费者 2000 条/s,则至少需要 5 个分区(且对应 5 个消费者实例)。
分区过多的代价
- Broker 端文件句柄、索引内存、请求对象成倍增加;
- Controller 故障转移时需要处理的分区状态变多,切换变慢;
- 生产端批次数增多,单批变小,压缩与吞吐效率下降;
- 未_leader 均衡或再分配时,迁移时间显著变长。
建议:起步用 6/12/24 这类预留量,宁可中期扩一次,也不要一开始上千分区。
7.2 副本因子与放置
生产建议:
replication.factor = 3(多数场景);min.insync.replicas = 2;acks=all的生产者配合,可容忍一台 Broker 故障且不丢已确认消息;- 副本应分布在不同机架/可用区(
broker.rack+replica.selector.class),防止机架级故障导致分区全丢。
权衡:RF=2 时若一台宕机,ISR 剩 1,min.insync.replicas=2 会导致写入不可用(这是保护一致性的正确行为);RF=3 则仍可写。
7.3 命令行管理回顾与补充
# 创建:显式指定分区与副本
kafka-topics.sh --create --topic payments \
--partitions 12 --replication-factor 3 \
--config retention.ms=604800000 \
--config min.insync.replicas=2 \
--config cleanup.policy=delete
# 查看某分区详细
kafka-topics.sh --describe --topic payments --partitions 0
# 修改 topic 级配置
kafka-configs.sh --bootstrap-server localhost:9092 \
--alter --entity-type topics --entity-name payments \
--set-config retention.ms=259200000
# 查看生效配置(含默认值)
kafka-configs.sh --bootstrap-server localhost:9092 \
--describe --entity-type topics --entity-name payments --all
7.4 用 AdminClient 管理主题
import org.apache.kafka.clients.admin.*;
import org.apache.kafka.common.config.ConfigResource;
import java.util.*;
import java.util.concurrent.ExecutionException;
public class TopicAdmin {
public static void main(String[] args) throws Exception {
Properties props = new Properties();
props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
try (AdminClient admin = AdminClient.create(props)) {
// 1. 创建主题
NewTopic newTopic = new NewTopic("payments", 12, (short) 3)
.configs(Map.of(
"retention.ms", "604800000",
"min.insync.replicas", "2"
));
admin.createTopics(List.of(newTopic)).all().get();
// 2. 查看描述
var desc = admin.describeTopics(List.of("payments"))
.allTopicNames().get().get("payments");
System.out.println("partitions=" + desc.partitions().size());
// 3. 增加分区
admin.createPartitions(Map.of("payments",
NewPartitions.increaseTo(24))).all().get();
// 4. 修改配置
ConfigResource res = new ConfigResource(
ConfigResource.Type.TOPIC, "payments");
admin.incrementalAlterConfigs(Map.of(res, List.of(
new AlterConfigOp(
new ConfigEntry("retention.ms", "259200000"),
AlterConfigOp.OpType.SET)
))).all().get();
// 5. 删除主题
// admin.deleteTopics(List.of("payments")).all().get();
}
}
}
AdminClient 是构建自服务管控台(申请 topic、查详情、改保留期)的基础。
7.5 增加分区的副作用
Kafka 只允许增加分区,不允许减少。增加时要注意:
- 新分区立即参与哈希计算,相同 key 的消息从此可能落进新分区,该 key 的新旧消息在全局上不再连续有序;
- 已有数据的分布不变,不会自动重分布(因此新分区是“空”的,短时间负载可能不均);
- 消费组会触发再平衡,把新分区分配给成员。
如果业务强依赖 key 顺序,两个选择:
- 一开始把分区规划到位,上线后不再扩;
- 升级时采用“新 topic + 双写/迁移”方案,消费端按新 key 顺序切换。
7.6 保留策略与日志压缩
| 场景 | 推荐 cleanup.policy | 说明 |
|---|---|---|
| 业务事件、日志 | delete | 按 retention.ms/retention.bytes 滚动删除 |
| 状态快照、配置变更、CDC 最新值 | compact | 每个 key 保留最新一条 |
| 想先压缩再过期 | compact,delete | 同时生效 |
日志压缩要点:
- key 为空的消息在 compact 主题中会被直接丢弃,必须带 key;
- 压缩是后台异步的,不是“写入即压缩”;
min.cleanable.dirty.ratio(默认 0.5)控制触发阈值:脏数据占比超过一半才清理;- 适合“当前状态”而非“事件历史”。
7.7 主题命名与治理
推荐 <域>.<数据集>.<事件类型> 这类规范:
order.domain.created.v1
user.profile.changed.v1
log.nginx.access.v2
治理要点:
- 命名带版本号,schema 破坏性变更走新版本主题;
- 禁止
auto.create.topics.enable=true(生产环境),防止拼错 topic 名悄悄建出新主题; - 用配额(quota)限制异常生产者,防止单个业务打挂集群;
- 建立 topic 清单与责任人登记,避免“僵尸主题”占用磁盘。
本章小结
- 分区数由目标吞吐和消费并行度共同决定,预留余量但不盲目多建;
- 生产标准:RF=3、min.insync.replicas=2、跨机架放置;
- 增分区会改变 key 哈希落点,对顺序敏感的业务要提前规划或走新主题迁移;
- compact 适合“每 key 最新状态”,delete 适合事件与日志;
- AdminClient 可以把主题管理沉淀为自服务能力。
思考题
- 一个 topic 从 6 分区扩到 12 分区后,为什么可能短时间出现新旧分区数据量严重不均?
- 为什么 Kafka 不支持减少分区?如果强行做,会遇到哪些一致性难题?
- 用户新增一个“查询用户当前等级”的需求,你会建 delete 还是 compact 主题?为什么?