KafkaNotes

第 07 章:主题与分区管理

zjc 于 2026-01-07 发布

这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 主题与分区是 Kafka 的“表结构”,一旦上线就很难随意修改。本章讲清楚分区数怎么定、副本因子怎么配、以及用 AdminClient 做自动化管理。

7.1 分区数怎么规划

分区数影响三个维度:

  1. 并行度:消费者实例上限 = 分区数;
  2. 吞吐:分区分布在多个 Broker 上,写入与读取并行;
  3. 资源与开销:每个分区都是一个目录、多组文件、多个副本、多个网络请求对象。分区不是越多越好。

经验公式

目标吞吐 / 单分区吞吐 = 所需分区数

例:目标 100 MB/s,压测单分区 10 MB/s -> 至少 10 个分区
   考虑峰值与增长,规划为 20-30 个

再检查消费端:如果下游每秒只能处理 1 万条、单消费者 2000 条/s,则至少需要 5 个分区(且对应 5 个消费者实例)。

分区过多的代价

建议:起步用 6/12/24 这类预留量,宁可中期扩一次,也不要一开始上千分区。

7.2 副本因子与放置

生产建议:

权衡: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 顺序,两个选择:

  1. 一开始把分区规划到位,上线后不再扩;
  2. 升级时采用“新 topic + 双写/迁移”方案,消费端按新 key 顺序切换。

7.6 保留策略与日志压缩

场景 推荐 cleanup.policy 说明
业务事件、日志 delete retention.ms/retention.bytes 滚动删除
状态快照、配置变更、CDC 最新值 compact 每个 key 保留最新一条
想先压缩再过期 compact,delete 同时生效

日志压缩要点:

7.7 主题命名与治理

推荐 <域>.<数据集>.<事件类型> 这类规范:

order.domain.created.v1
user.profile.changed.v1
log.nginx.access.v2

治理要点:

本章小结

思考题

  1. 一个 topic 从 6 分区扩到 12 分区后,为什么可能短时间出现新旧分区数据量严重不均?
  2. 为什么 Kafka 不支持减少分区?如果强行做,会遇到哪些一致性难题?
  3. 用户新增一个“查询用户当前等级”的需求,你会建 delete 还是 compact 主题?为什么?