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 Hash | hash(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.partitions | 1 | 按吞吐计算 | 分区数 |
replication.factor | 1 | 3(生产环境) | 副本数 |
min.insync.replicas | 1 | 2 | 最小同步副本数 |
retention.ms | 7 天 | 按需 | 消息保留时间 |
retention.bytes | -1(无限制) | 磁盘限制 | 按大小保留 |
segment.ms | 7 天 | 1 天 | 滚动新日志段的时间 |
cleanup.policy | delete | delete/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 | 磁盘爆满是迟早的事 | 根据业务设置保留策略 |
| 盲目增加分区 | 元数据膨胀、启动变慢 | 评估实际吞吐需求 |
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。