04. Topic 设计与分区策略

Kafka Topic 分区数设计、键分区顺序保证、Topic 命名规范与最佳实践,掌握高吞吐场景下的 Topic 规划方法论。

1. Topic 核心概念

1.1 Topic 是什么

Topic 是 Kafka 中消息的逻辑分类。Producer 将消息发送到特定 Topic,Consumer 订阅 Topic 获取消息。Topic 本身不存储数据,真正的数据存储在 Topic 的**分区(Partition)**中。

Topic: orders
  ├── Partition 0 → [M1, M2, M5, M9,  ...]  → Broker 101
  ├── Partition 1 → [M3, M6, M7, M10, ...]  → Broker 102
  └── Partition 2 → [M4, M8, M11, M12, ...] → Broker 103

每个 Partition 是一个有序的、不可变的消息序列
Partition 内的消息顺序保证,Partition 间无序

1.2 分区的作用

作用说明
并行度分区数决定消费者最大并行度
负载均衡分区均匀分布到各 Broker,实现负载均衡
顺序保证同一分区内消息严格有序
水平扩展增加分区即可提升 Topic 吞吐量

2. 分区数设计

2.1 分区数计算公式

分区数 = max(
    预期吞吐量 / 单分区吞吐量,
    消费者数(期望并行度)
)

单分区吞吐量上限(典型值):
  - 生产者写入:~10 MB/s(取决于磁盘 IOPS 和网络带宽)
  - 消费者读取:~10 MB/s
  
示例:
  预期吞吐:100 MB/s
  单分区吞吐:10 MB/s
  期望并行:12 个消费者
  
  分区数 = max(100/10, 12) = max(10, 12) = 12

2.2 分区数过多/过少的风险

问题原因后果
分区过少消费者数 > 分区数部分消费者空闲,资源浪费
分区过多分区数 » 消费者数文件句柄耗尽、启动慢、元数据膨胀
分区过多每个分区一个目录 + 日志段磁盘随机读写增加

经验法则:单 Broker 承载分区数不超过 2000-4000,单集群不超过 20 万分区。

2.3 调整分区

# 增加分区(只能增不能减)
kafka-topics.sh --bootstrap-server kafka:9092 \
    --alter --topic orders \
    --partitions 12

# ⚠️ 注意:分区增加后,key 分区的顺序保证可能被打破!
# 原因:新的消息可能路由到新增的分区,导致同一 key 的消息分布在不同分区

3. 分区策略

3.1 DefaultPartitioner(默认策略)

// Kafka 2.4+ 默认分区策略(Sticky Partitioner)
public int partition(String topic, Object key, byte[] keyBytes,
                     Object value, byte[] valueBytes, Cluster cluster) {
    if (keyBytes == null) {
        // 无 key:Round-Robin 到各分区(批次级别)
        return stickyPartitionCache.partition(topic, cluster);
    }
    // 有 key:hash(key) % numPartitions
    return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;
}

3.2 分区策略对比

策略方式顺序保证适用场景
RoundRobin轮询分配无顺序要求,追求最大吞吐
Key Hashhash(key) % n同一 key 有序需要按 key 聚合的统计
Custom自定义逻辑自定义特殊路由需求(如按地域)
Sticky批次分配到同一分区减少延迟,提高吞吐(默认)

3.3 自定义分区器

// 按地域分区,同地域消息进入同一分区
public class RegionPartitioner implements Partitioner {
    @Override
    public int partition(String topic, Object key, byte[] keyBytes,
                         Object value, byte[] valueBytes, Cluster cluster) {
        String region = extractRegion(value);
        int numPartitions = cluster.partitionCountForTopic(topic);
        return Math.abs(region.hashCode()) % numPartitions;
    }
    
    private String extractRegion(Object value) {
        // 从消息体解析地域
        return ((Order) value).getRegion();
    }
    
    @Override public void configure(Map<String, ?> configs) {}
    @Override public void close() {}
}

// 配置 Producer
props.put("partitioner.class", "com.example.RegionPartitioner");

4. 顺序保证

4.1 Kafka 顺序保证范围

Kafka 的顺序保证:
  ✅ 同一 Partition 内,消息按写入顺序消费
  ✅ 同一 Producer 的同一 Partition 内,消息顺序保证
  ❌ 不同 Partition 间,消息无顺序保证
  ❌ 重试可能导致同 Partition 内乱序(max.in.flight > 1 + 重试)

4.2 实现全局有序的方案

方案 1:单分区 Topic(牺牲吞吐)
  Topic: orders-single-partition (partition=1)
  所有消息进入同一分区 → 全局有序
  适用:低频但要求严格顺序的场景(如账户流水)

方案 2:Key 分区(局部有序)
  Producer:
    send("orders", userId, message);
    
  同一 userId 的 hash 结果相同 → 同一分区
  → 同一用户的订单有序
  
  适用:用户级别的顺序要求(如用户订单、用户消息)

4.3 防止重试导致的乱序

// 方案:max.in.flight.requests.per.connection = 1
props.put("max.in.flight.requests.per.connection", 1);

// 缺点:吞吐降低(每次只能发一个请求,需等待确认后才能发下一个)
// Kafka 2.5+ 启用幂等生产者后,可安全设置为 5
props.put("enable.idempotence", true);  // 默认开启
// 开启幂等后,Kafka 自动保证 per-partition 顺序,即使 inflight > 1

5. Topic 命名规范

5.1 命名建议

格式:<domain>.<event>.<version>

示例:
  commerce.order.created.v1
  commerce.order.cancelled.v1
  user.profile.updated.v1
  inventory.stock.reserved.v1
  payment.transaction.completed.v1

好处:
  1. 清晰表达业务领域和事件类型
  2. 版本号便于演进(v1 → v2 不破坏现有消费者)
  3. ACL 权限可按 domain 前缀批量配置

5.2 Topic 配置参数

参数默认值建议说明
num.partitions1按吞吐计算分区数
replication.factor13(生产环境)副本数
min.insync.replicas12最小同步副本数
retention.ms7 天按需消息保留时间
retention.bytes-1(无限制)磁盘限制按大小保留
segment.ms7 天1 天滚动新日志段的时间
cleanup.policydeletedelete/compact清理策略

5.3 日志压缩(Log Compaction)

适用场景:需要保留 key 的最新值(如配置更新、用户资料)

普通删除策略:
  -> 时间/大小到了 → 删除整个日志段

压缩策略:
  -> 只保留每个 key 的最新值,旧值被清理
  
  时间线:
    t1: K1=V1, K2=V1
    t2: K1=V2      ← K1=V1 被清理
    t3: K2=V2, K3=V1
    t4: K1=V3      ← K1=V2 被清理
    
  最终保留:K1=V3, K2=V2, K3=V1
# 创建 compact topic
kafka-topics.sh --create --topic user-config \
    --config cleanup.policy=compact \
    --config min.cleanable.dirty.ratio=0.5 \
    --config delete.retention.ms=100

6. Topic 设计最佳实践

6.1 设计检查清单

□ 吞吐评估:预计消息量、消息大小、峰值 QPS
□ 分区计算:分区数 = max(吞吐/单分区上限, 消费并行度)
□ 副本设置:生产环境 replication.factor >= 3
□ 顺序需求:是否需要 key 分区?还是全局单分区?
□ 保留策略:按时间保留还是压缩保留?
□ 命名规范:domain.event.version 格式
□ 扩展预留:初期分区数应大于当前消费者数(预留扩展空间)

6.2 常见错误

错误后果解决
分区数 = 消费者数新增消费者无法加入分区数 = 消费者数 × 2~3
忽略 key 分区同一用户订单乱序使用 userId 作为 key
不设置 retention磁盘爆满是迟早的事根据业务设置保留策略
盲目增加分区元数据膨胀、启动变慢评估实际吞吐需求

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「kafka」更多文章

  1. 事件驱动架构:Event Sourcing、CQRS 与 Saga 模式
  2. Kafka 运维监控与故障恢复:JMX 指标、Lag 监控与分区重分配
  3. Kafka 详解:分布式日志系统、ISR 与一致性保证