主题
第 6 章 Topic 设计与分区策略:从「能跑」到「跑得稳」
目标读者:会用
kafka-topics.sh --create拍脑袋建 Topic、副本数永远填 1、分区数永远填 1 或者上百个,然后线上要么扛不住要么 Rebalance 抖死的同学。学完你会:拿到一个新业务,能在 5 分钟内说清楚「这个 Topic 应该开多少分区、多少副本、Key 怎么定、
min.insync.replicas填几、未来扩容怎么办」,并且讲得出每一项背后的依据。
0. 导读:分区数和副本数,是 Kafka 工程师的「灵魂两问」
绝大多数线上 Kafka 事故的根因,可以归结成下面这几条之一:
- 分区数选少了 —— 后端再多消费者也吃不动,
records-lag-max一路涨;想加分区,但 Key Hash 分布会变,下游所有「按 Key 聚合」的逻辑全乱。 - 分区数选多了 —— 单 Broker 上几万个 Partition,Controller 元数据爆炸、
OpenFileDescriptors顶到上限、宕机后 Recovery 几十分钟,新 Leader 怎么也选不出来。 - 副本数填了 1 —— 一台 Broker 宕机,整个 Topic 直接
OFFLINE,业务方电话打爆。 acks=1+min.insync.replicas=1—— 看着「保留 3 副本」,其实写下去只要 Leader 自己确认就 ack,Leader 一挂数据丢了一片。- Key 用了
userId但热门用户占 80% 流量 —— 单个分区直接被打爆,其它分区闲着发呆。 - 「先小心试运行」副本数 1,后面想加副本 ——
kafka-reassign-partitions.sh触发数据复制,整个集群网卡跑满,业务延迟飙到秒级。
这些问题没有一条是「Kafka 不好用」,全部是 Topic 设计阶段没想清楚。本章就是要让你在拍下 --partitions N --replication-factor M 之前,心里有数。
1. 一个生活类比:高速公路收费站
把一个 Topic 想象成一条高速公路:
- Partition = 一条车道。车道越多,单位时间能放过去的车(消息)就越多。
- Replica = 同一段路的备用车道(在不同地块),主车道塌方了备用车道顶上。
- Producer = 入口收费站,按车牌(Key)分配车道。
- Consumer Group = 出口,每个出口同时刻只能站一个收费员,但一个收费员可以同时管多个车道。
- Key = 车牌(决定走哪个车道)。
- Offset = 车道上的车序号。
这个类比能直接解释 4 条最重要的规则:
| 规则 | 对应的高速公路常识 |
|---|---|
| 单分区内消息严格有序 | 一条车道上的车前后顺序固定 |
| 跨分区无序 | 三条车道的车并排开,没人能保证整体顺序 |
| Key 相同必然落同一分区 | 同一个车牌总是走同一条车道 |
| 加分区会导致 Key 重分布 | 多开了一条车道,原车道上一半的车被分流走了 |
记住「车道」这张图,下面所有公式和取舍都能秒懂。
2. 分区数(partitions)怎么选?
2.1 经验公式
业界最常引用的就是这条:
partitions = max( T / P_p , T / C_c )| 符号 | 含义 | 单位 | 怎么测 |
|---|---|---|---|
T | 业务目标吞吐 | MB/s 或 msg/s | 业务方给 |
P_p | 单分区生产吞吐 | MB/s 或 msg/s | 用 kafka-producer-perf-test.sh 在 1 分区 Topic 上压 |
C_c | 单消费者吞吐 | MB/s 或 msg/s | 业务消费逻辑跑一次实测 |
含义:分区数同时决定「写入并行度」和「消费并行度」。
- 想要
T的写入吞吐,至少要T / P_p条「车道」让 Producer 写。 - 想要
T的消费吞吐,至少要T / C_c条「车道」让 Consumer 拉。 - 取大的那个,因为系统瓶颈以「最慢的那一段」为准。
2.2 一个具体例子
业务给的指标:
- 目标吞吐:100 MB/s
- 单分区生产实测:40 MB/s(普通 SATA SSD、
acks=all、compression=lz4、批次 64KB) - 单消费者实测:20 MB/s(业务里要把 JSON 反序列化、写 ClickHouse)
代入公式:
T / P_p = 100 / 40 = 2.5 → 取 3
T / C_c = 100 / 20 = 5 → 取 5
partitions = max(3, 5) = 5实际我们一般会再 乘以 1.5 ~ 2 的余量,因为:
- 业务高峰可能是平均的 2 倍。
- 未来一年业务可能翻倍,加分区有破坏性(见 §6),不如一开始多开。
- 消费者数量上限就是分区数,组里多挂一个也不会浪费。
最终选 8 ~ 12 个分区,是这个例子里比较稳的答案。
2.3 为什么不能「无脑开 100 个分区」?
「反正分区越多越好,开 200 个总错不了。」 —— 这句话在 ZK 时代会让你的 Controller 直接宕机。
每多一个分区,集群要承担的成本:
| 成本项 | 说明 |
|---|---|
| 元数据条目 | 每个分区的 Leader、ISR、Epoch 都要在 Controller 内存里存一份;ZK 时代单 Controller 5 万分区基本就是上限,KRaft 提到 200 万但元数据传播延迟仍然增长。 |
| 文件句柄 | 每个分区至少 3 个文件(.log / .index / .timeindex),开 1 万分区就是 3 万 fd,外加 Producer/Consumer 的 socket,ulimit -n 不够直接打开失败。 |
| PageCache 浪费 | 每个 Segment 起头几 KB 都要进 PageCache,分区越多冷热数据混杂越严重,PageCache 命中率下降。 |
| Replication 线程压力 | Follower 要拉每个 Partition 的数据,分区多 = Fetch 请求批数多,CPU 飙升。 |
| Recovery 时间 | Broker 重启时要重放每个分区的最后一段日志,10 万分区可能要 30 分钟才上线。 |
| End-to-End 延迟 | Producer 端一次 flush() 要等所有命中的分区 ack,分区越多每个分区的 batch 就越小,linger 凑不满,延迟反而大。 |
经验值:
| 集群规模 | 单 Broker 推荐 Partition 上限 |
|---|---|
| 3 ~ 5 节点测试集群 | ≤ 2,000 |
| 中型生产集群(10 节点左右) | ≤ 4,000 |
| 大型 KRaft 集群(30+ 节点) | ≤ 8,000,全集群不超 200,000 |
2.4 单 Topic 多大算合理?
单 Topic 分区数推荐区间:6 ~ 60。
- 下限 6:让消费者组至少能扩到 6 个实例,留扩容空间。
- 上限 60:再多就要看是不是真的需要——经常是 Key 设计错了,应该改 Key 而不是加分区。
特殊场景:
- 海量小 Topic(每个 Topic 流量很小,但 Topic 数量上千):每个 Topic 1 ~ 3 个分区即可,重点控总分区数。
- 超大 Topic(单业务流量 GB/s 级):可能 100+ 分区,但要配合 broker 数量同步上调,并且强烈建议 KRaft。
__consumer_offsets、__transaction_state等内部 Topic:默认 50 分区,不要乱改。
3. 副本数(replication.factor)怎么选?
3.1 三档基准
| RF | 适用场景 | 风险 |
|---|---|---|
1 | 本地 Demo / 个人学习,绝不允许进生产 | Broker 宕机即数据丢失 |
2 | 资源紧张的非核心日志(埋点、监控) | 一台 Broker 挂掉立刻只剩 1 副本,等于 RF=1,再坏一个就丢数据 |
3 | 生产默认值,强烈推荐 | 容忍 1 台 Broker 宕机仍有 2 副本,配合 min.insync.replicas=2 满足 CAP 中的 CP |
5 | 跨机房 / 跨可用区的金融级业务 | 网络复制开销大,但能容忍 2 台同时挂 |
3.2 「3 副本 + min.insync.replicas=2 + acks=all」黄金三角
这套组合是 Kafka 在「可用性」和「一致性」之间最经典的妥协:
properties
# Topic 级
replication.factor=3
min.insync.replicas=2
# Producer 级
acks=all
enable.idempotence=true语义解读:
- 写入时,Leader + 至少 1 个 Follower 都把消息写进各自的日志,Producer 才收到 ack。
- 读取时,Consumer 只能读到「HW(High Watermark)」以下的消息,HW 由 ISR 中所有副本的 LEO 取 min 得到。
- 如果 ISR 缩到只剩 1 个(满足不了
min.insync.replicas=2),Producer 端会收到NotEnoughReplicasException,拒绝写入而不是默默丢数据。
为什么不是 min.insync.replicas=3?
如果 RF=3、min.insync.replicas=3,那么任何一台 Broker 短暂抖动都会把 ISR 缩到 2,Producer 立刻被拒。可用性会非常差。min.insync.replicas=2 在 RF=3 下兼顾了「至少 2 副本落盘」和「容忍 1 副本临时掉队」。
3.3 副本数的成本账
每多一个副本,等于:
- 磁盘空间 ×N
- 网络入口流量 ×N(Follower 要拉数据)
- 网络出口流量 Leader ×(N-1)
所以 RF=5 的代价大约是 RF=3 的 1.67 倍存储 + 1.67 倍网络。除非业务真的承受不了「一次断电丢一段消息」的风险,否则 RF=3 就够了。
3.4 跨机房的特殊处理
如果集群横跨 2 ~ 3 个可用区(AZ),单纯设 RF=3 不够,因为 Kafka 默认不知道副本分布在哪个机房,可能 3 个副本全分到同一 AZ。需要:
bash
# Broker 端配置(每台 Broker 标注自己的机架/AZ)
broker.rack=az-a # 或 az-b, az-c
# Topic 创建时启用机架感知(默认开启,只要 broker.rack 设了就生效)
kafka-topics.sh --create \
--topic learn.06.cross-az \
--partitions 6 \
--replication-factor 3 \
--bootstrap-server kafka-1:9092Kafka 会自动让 3 个副本分布到 3 个不同的 broker.rack,这样任何一个机房整体宕掉,集群仍能写入。
4. Key 的选择:决定顺序、决定热点
4.1 Key 的两个作用
作用 1:决定分区
python
# Producer 端伪流程(confluent-kafka 默认 partitioner)
if key is None:
partition = sticky_random_partition() # 没有 Key 时走 Sticky Partitioner
else:
partition = murmur2(key) % num_partitions # Murmur2 哈希作用 2:决定顺序
同一个 Key 的消息一定会落到同一个分区,所以同一个 Key 的消息一定有序。 不同 Key 的消息没有任何顺序保证。
4.2 选 Key 的三条原则
原则 1:业务必须按 Key 顺序处理时,必须设 Key
典型反例:订单状态机消息 {order_id, status}。
- 如果不设 Key,订单 A 的「待支付」「已支付」「已发货」可能分到 3 个不同分区,被 3 个 Consumer 并发处理,下游可能先看到「已发货」再看到「已支付」,状态机直接错乱。
- 正确做法:
key = order_id,同一个订单的所有事件强制走同一分区,串行消费,状态机天然正确。
原则 2:Key 的基数(cardinality)要远大于分区数
如果只有 5 个 Key、却开了 12 个分区,Murmur2 哈希也救不了,至少 7 个分区永远空着。
- 推荐:基数 ≥ 分区数 × 100。
- 例:分区 12 个,Key 至少 1200 种以上才能保证基本均衡。
原则 3:Key 的分布要尽量均匀,避免热点
Key 选错的最常见反例:用
tenantId当 Key,结果某个大客户占了 80% 流量,整个 Topic 的吞吐被一个分区拖垮。
排查方法:
python
# 用 hash 模拟分布(详见 06_topic_design/code/key_distribution.py)
from collections import Counter
counts = Counter(murmur2(k) % 12 for k in real_keys)
print(counts)
# 健康:每个分区占比在 1/12 ± 30% 范围内
# 不健康:某分区占比 >2/12修复方案:
- 把 Key 加随机后缀打散:
key = f"{tenantId}_{random.randint(0,9)}",但会破坏顺序,慎用。 - 业务允许的话,改 Key 维度,比如改用
order_id而不是tenantId。 - 真的要保留 tenant 顺序,又要扛热点:开专用 Topic 给大租户,业务侧路由。
4.3 Murmur2 + 取模演示
Murmur2 是 Kafka 默认 Partitioner 的哈希算法,它的特点是均匀、非加密、非常快。
简化算法(实际见 org.apache.kafka.common.utils.Utils#murmur2):
python
def murmur2(data: bytes) -> int:
seed = 0x9747b28c
m = 0x5bd1e995
r = 24
length = len(data)
h = seed ^ length
i = 0
while length >= 4:
k = (data[i] & 0xff) | ((data[i+1] & 0xff) << 8) \
| ((data[i+2] & 0xff) << 16) | ((data[i+3] & 0xff) << 24)
k = (k * m) & 0xffffffff
k ^= (k >> r)
k = (k * m) & 0xffffffff
h = (h * m) & 0xffffffff
h ^= k
i += 4
length -= 4
if length == 3: h ^= (data[i+2] & 0xff) << 16
if length >= 2: h ^= (data[i+1] & 0xff) << 8
if length >= 1:
h ^= (data[i] & 0xff)
h = (h * m) & 0xffffffff
h ^= (h >> 13)
h = (h * m) & 0xffffffff
h ^= (h >> 15)
return h & 0x7fffffff # toPositive最后落点:
partition = murmur2(key_bytes) & 0x7fffffff % num_partitions关键观察:% num_partitions 这一步——只要 num_partitions 变了,几乎所有 Key 的落点都会变。这就是 §6 要展开讲的「重建 Topic vs 加分区」难题。
4.4 没有 Key 怎么办?Sticky Partitioner
python
producer.send("topic", value=b"hello") # 没传 keyKafka 2.4+ 默认走 Sticky Partitioner:
- 选定一个分区,整个 batch 都发同一个分区。
- 等 batch 被发出去(达到
linger.ms或batch.size)之后,再换一个分区。
老版本是「Round-Robin」每条换分区,会让每个分区只攒到一条消息,batch 永远凑不满。Sticky 之后吞吐能提升 2 ~ 3 倍。
没 Key 的副作用:消息分布是「按时间分块」的,不是「按 Key 哈希分均匀」的。如果业务真的要均匀分布且不需要顺序,可以加一个随机 Key:
python
producer.send("topic", key=str(uuid.uuid4()).encode(), value=value)5. min.insync.replicas 与 acks 的配合
5.1 acks 的三档语义
acks | Producer 何时收到 ack | 一致性 | 性能 |
|---|---|---|---|
0 | 发出去就算成功,不等任何 ack | 可能丢消息(Broker 没收到也不知道) | 最快 |
1 | Leader 写入本地日志(不一定 fsync)就 ack | Leader 突然挂掉、还没复制给 Follower → 丢消息 | 中等 |
all(= -1) | ISR 中所有副本都已经写入本地日志才 ack | 配合 min.insync.replicas 才有意义 | 最慢 |
5.2 acks=all 不等于「写到所有副本」
这是新手最常见误区:
「我设了
acks=all,3 副本都写进去才 ack,万无一失了吧。」
错。acks=all 是「ISR 里所有副本都 ack」,不是「所有副本」。如果 ISR 因为某个 Follower 掉队缩成只剩 1 个(也就是只剩 Leader),那么 acks=all 等同于 acks=1。
所以必须配合:
properties
min.insync.replicas=2 # ISR 至少要 2 个,才允许写含义:ISR 缩到 1 时,Producer 直接收到 NotEnoughReplicasException,拒写而不是默默降级。这就把「丢消息」变成了「业务可观测的错误」,业务可以重试或报警。
5.3 黄金三角的故障行为表
假设 RF=3, min.insync.replicas=2, acks=all:
| 场景 | ISR 状态 | 写入结果 |
|---|---|---|
| 一切正常 | 成功 | |
| F1 掉队被踢出 | 成功(ISR=2 ≥ 2) | |
| F1 和 F2 都掉队 | 拒写(NotEnoughReplicasException) | |
| Leader 宕机 | 重新选主,可能 | 短暂不可用,几秒后恢复 |
| 全部 Broker 都挂 | — | 不可用 |
5.4 Producer 端必带的兜底配置
properties
# 强一致 Producer 推荐配置(business critical 场景)
acks=all
enable.idempotence=true # 自动开 max.in.flight=5、retries=Integer.MAX、acks=all
max.in.flight.requests.per.connection=5 # 必须 ≤ 5 才能保证顺序
delivery.timeout.ms=120000 # 超过这个时间还没 ack 就报错
retries=2147483647 # 给重试一个上限「无限大」,由 delivery.timeout 控制总时长6. 单分区顺序 vs 跨分区无序:「车道」边界
6.1 三句话讲清楚
- 单分区内:消息 100% 严格按写入顺序,Consumer poll 出来也 100% 按顺序。
- 跨分区:Kafka 不提供任何顺序保证。
- 跨 Topic:完全没有顺序概念。
6.2 一个常见误区
「我业务里订单状态变更走 Topic A、订单支付走 Topic B,怎么保证 A 先于 B?」
答案:Kafka 帮不了你。要么合并到一个 Topic(同 Key),要么在业务层加版本号(state version、occurred_at)排序。
6.3 跨分区无序的视觉化
分区 0: m1 ─ m4 ─ m7 ─ m10
分区 1: m2 ─ m5 ─ m8 ─ m11
分区 2: m3 ─ m6 ─ m9 ─ m12
Consumer poll 出来的顺序可能是:
m1, m2, m4, m3, m5, m6, m7, m8, m10, m9, m11, m12车道类比:3 条车道并行开,外面的人记录通过收费亭的车牌号,顺序是随机的。
6.4 怎么取「全局有序」?
只有一个办法:单分区 Topic(partitions=1)。
代价:
- 吞吐被锁死在单分区上限(可能 ≤ 50 MB/s)。
- 消费者组里只能有 1 个有效 Consumer。
- 未来扩容只能加机器,不能加并行度。
所以 Kafka 的设计哲学是:抛弃全局有序,换取无限水平扩展。如果你真的需要全局有序,先问自己「是不是其实可以拆成多个独立的 key 序」。
7. 何时要「重建 Topic」而不是「加分区」?
7.1 加分区会破坏什么?
回到 §4.3 的公式:
partition = murmur2(key) % N把 N 从 6 改成 8:同一个 Key 的落点变化概率是 (N_new - N_old) / N_new = 2/8 = 25%。具体来说:
| Key | murmur2 % 6 | murmur2 % 8 | 落点是否变化 |
|---|---|---|---|
order_001 (hash=12345678) | 12345678 % 6 = 4 | 12345678 % 8 = 6 | 变了 |
order_002 (hash=98765432) | 98765432 % 6 = 2 | 98765432 % 8 = 0 | 变了 |
order_003 (hash=11111111) | 11111111 % 6 = 5 | 11111111 % 8 = 7 | 变了 |
大部分 Key 都会换分区。
7.2 谁会受伤?
- 消费者侧的「按 Key 聚合」状态
- 例:Kafka Streams 的
aggregate(key)把order_001的状态存在 partition 4 对应的 State Store 里。加分区后order_001的新消息去了 partition 6,但旧状态还在 partition 4 的 Store,状态对不上号。
- 例:Kafka Streams 的
- 顺序保证被打破
- 加分区瞬间,
order_001的新消息进了 partition 6,但 partition 4 上还有一条没消费完的旧消息。Consumer A(消费 P4)和 Consumer B(消费 P6)并发处理,旧消息可能在新消息之后才被处理。
- 加分区瞬间,
- 下游去重 / 幂等表的依赖
- 如果下游用
(partition, offset)当幂等键,加分区后这个键的语义就变了。
- 如果下游用
7.3 决策树
要不要加分区?
├─ Topic 没 Key(消息均匀洒下去) → 安全,可加
├─ Topic 有 Key 但下游不依赖顺序 → 可加,但要通知下游
├─ Topic 有 Key 且下游严重依赖顺序 → 必须「重建 Topic」
└─ 是 __consumer_offsets 等内部 Topic → 永远不要碰7.4 重建 Topic 的标准流程
bash
# 1. 创建新 Topic(更大的分区数)
kafka-topics.sh --create --bootstrap-server kafka:9092 \
--topic learn.06.orders.v2 \
--partitions 32 --replication-factor 3 \
--config min.insync.replicas=2
# 2. 启动 MirrorMaker 2 / 自研 bridge,把旧 Topic 的数据按 Key 重哈希复制到新 Topic
# (注意:这里是「重新走一次 Producer」,Key 会按新分区数重新哈希到位)
# 3. 等下游全部切到新 Topic 后,归档旧 Topic
kafka-topics.sh --delete --topic learn.06.orders.v1 --bootstrap-server kafka:9092业界俗称「Shadow Topic 双写切换」。
8. Topic 命名规范
一个生产集群里几百上千个 Topic,命名混乱会让运维抓狂。推荐规范:
<env>.<domain>.<entity>.<event>[.v<version>]| 字段 | 说明 | 例 |
|---|---|---|
env | 环境前缀(避免跨环境串) | prod / staging / dev |
domain | 业务域 | order / payment / user |
entity | 数据实体 | purchase / wallet |
event | 事件类型(用过去式) | created / paid / refunded |
version | 版本(重建 Topic 时 +1) | v1 / v2 |
例:
prod.order.purchase.created.v1prod.order.purchase.paid.v2dev.user.profile.updated.v1
教程统一用 learn.<chapter>.<scene> 是为了简短,但生产建议用上面的完整方案。
禁忌:
- 用中文 / 大写字母 / 下划线开头(
__是 Kafka 保留前缀,对应内部 Topic)。 - Topic 名超过 249 字符(Kafka 硬限制)。
- 把「下游消费者名」写进 Topic 名(如
order.for.risk.team),破坏 Pub/Sub 解耦。
9. 实战:建一个生产级 Topic
9.1 命令行
bash
kafka-topics.sh --bootstrap-server kafka:9092 --create \
--topic prod.order.purchase.created.v1 \
--partitions 12 \
--replication-factor 3 \
--config min.insync.replicas=2 \
--config retention.ms=604800000 \
--config segment.bytes=536870912 \
--config compression.type=producer \
--config max.message.bytes=1048576逐项解读:
| 配置 | 值 | 说明 |
|---|---|---|
partitions | 12 | 业务峰值 200MB/s ÷ 单分区 40MB/s ≈ 5,留 2 倍余量 |
replication.factor | 3 | 黄金三角 |
min.insync.replicas | 2 | 黄金三角 |
retention.ms | 604800000 | 7 天 |
segment.bytes | 512MB | 控制单 Segment 大小,避免太多小文件 |
compression.type | producer | 透传 Producer 端压缩,不做解压再压 |
max.message.bytes | 1MB | 防止巨型消息撑爆 PageCache |
9.2 模拟输出
Created topic prod.order.purchase.created.v1.9.3 校验
bash
kafka-topics.sh --bootstrap-server kafka:9092 \
--describe --topic prod.order.purchase.created.v1模拟输出:
Topic: prod.order.purchase.created.v1 TopicId: ABCDEFGHIJK PartitionCount: 12 ReplicationFactor: 3
Configs: min.insync.replicas=2,retention.ms=604800000,segment.bytes=536870912,compression.type=producer,max.message.bytes=1048576
Topic: prod.order.purchase.created.v1 Partition: 0 Leader: 1 Replicas: 1,2,3 Isr: 1,2,3
Topic: prod.order.purchase.created.v1 Partition: 1 Leader: 2 Replicas: 2,3,1 Isr: 2,3,1
Topic: prod.order.purchase.created.v1 Partition: 2 Leader: 3 Replicas: 3,1,2 Isr: 3,1,2
...注意 Replicas 列表的第一个元素就是「Preferred Leader」——优先 Leader。理想情况下 Leader 应该等于 Preferred Leader,分布均匀。如果不均,可用 kafka-leader-election.sh 重新选主。
10. 与 RabbitMQ / RocketMQ / Pulsar 的对比
| 维度 | Kafka | RabbitMQ | RocketMQ | Pulsar |
|---|---|---|---|---|
| 分区概念 | Topic-Partition 显式 | Queue 不分区,靠多 Queue + Exchange | MessageQueue(类似 Kafka Partition) | Topic-Partition + Bookkeeper Ledger |
| 顺序保证 | 单分区有序 | 单 Queue 有序 | 单 MessageQueue 有序 | 单分区有序 |
| 副本 | Leader-Follower (ISR) | 镜像队列(HA Mirror) | Master-Slave 同步 / 异步 | Bookkeeper 多副本 |
| 副本与存储分离 | 否(Broker 兼存计算) | 否 | 否 | 是(Broker 无状态 + Bookkeeper 存储) |
| 加分区代价 | Key Hash 重分布,破坏性大 | 增 Queue 不影响旧 Queue | 类似 Kafka,破坏性 | 较小(Bookkeeper 层透明扩 Bookie) |
| Key 路由 | Murmur2 哈希到分区 | 由 Routing Key + Exchange 类型决定 | Key/Tag 路由 | Key 哈希到分区 |
| 推荐 RF | 3 | 3 节点镜像 | 2 主 2 从 | 2 写 2 副本(quorum) |
一句话差异:Kafka 的「Topic 设计」是一锤子买卖,分区数和 Key 一旦定下来,改起来代价很大;RabbitMQ 因为 Queue 不分区,加 Queue 几乎无代价,但顺序粒度也更细;Pulsar 因为存算分离,扩 Bookie 几乎免费,但运维复杂度更高。
11. 小结
Topic 设计四件套
├─ partitions
│ ├─ 公式:max(T/P_p, T/C_c) × 1.5~2 余量
│ ├─ 上下限:单 Topic 6~60,单 Broker ≤ 4000
│ └─ 加分区破坏性:Key Hash 重分布
├─ replication.factor
│ ├─ 推荐 3
│ ├─ 跨机房用 broker.rack 机架感知
│ └─ 成本 = N × 存储 + N × 网络
├─ min.insync.replicas
│ ├─ 推荐 2(配 RF=3)
│ └─ 必须配 acks=all 才有意义
└─ Key
├─ 决定分区 + 决定顺序
├─ 基数 ≥ 分区数 × 100
└─ 防热点:均匀分布、必要时打散12. 面试高频题(6 题)
Q1:分区数怎么选?给一个公式和一个例子。
考察点:吞吐建模、容量规划。
答案:
- 公式:
partitions = max(T/P_p, T/C_c) × 1.5~2 余量,其中T是业务目标吞吐,P_p是单分区生产吞吐(实测),C_c是单消费者吞吐(实测)。 - 例:目标 100MB/s,单分区写 40MB/s,单消费者 20MB/s →
max(100/40, 100/20) = max(3, 5) = 5,再 ×2 余量 = 10 个分区。 - 上限考量:单 Broker 一般不超 4000 分区,否则 Controller 元数据、文件句柄、Recovery 时间都会炸。
- 不能无脑大:分区越多,每个 batch 越凑不满,端到端延迟反而上升;同时加分区会让 Key Hash 重分布,破坏顺序。
- 加分项:提到 Sticky Partitioner(无 Key 时)让 batch 集中、不至于分得太散;提到 Kafka 4.x KRaft 下分区上限提高到几十万但仍要看具体硬件。
Q2:副本数 3 + min.insync.replicas=2 + acks=all 这套组合,能保证什么、不能保证什么?
考察点:副本机制、可靠性语义。
答案:
- 能保证:消息一旦 ack,至少在 2 个 Broker 的本地日志里。容忍 1 台 Broker 同时宕机不丢数据。
- 不能保证:ISR 缩到 1 个时 Producer 拒写(
NotEnoughReplicasException),所以业务方要做好重试和报警。 - acks=all 的细节:是「ISR 中所有副本」ack,不是「所有副本」。所以必须配
min.insync.replicas才安全。 - 故障行为:
- 1 台 Broker 挂:写入仍然成功(ISR 还剩 2)。
- 2 台 Broker 挂(且 Leader 是其中之一):触发 Leader 选举,几秒不可用,但不丢数据。
- 全集群崩:完全不可用,但只要重启起来数据还在。
- 加分项:提到
unclean.leader.election.enable=false(默认 false)保证不会让 Out-of-Sync 副本当 Leader 导致 ack 过的消息被截断;提到enable.idempotence=true解决 retries 引起的重复。
Q3:什么情况下消息能保证顺序?跨分区怎么办?
考察点:顺序语义、Key 设计。
答案:
- 单分区顺序:同一个分区的消息,写入顺序 == 存储顺序 == 消费顺序。
- 跨分区无序:Kafka 不提供任何跨分区顺序保证。
- 怎么落到同分区:Producer 设 Key,Kafka 用
murmur2(key) % num_partitions决定分区。Key 相同 → 分区相同 → 顺序保证。 - 业务层面的解决方案:
- 同一聚合根(订单号、用户 id)用同一 Key,让相关事件同分区。
- 必要时按版本号 / 时间戳在消费者端排序。
- 极端情况下用单分区 Topic 实现全局有序(牺牲吞吐)。
- 顺序的隐藏破坏者:
max.in.flight.requests.per.connection > 1+ 关掉幂等:retry 会乱序。开了幂等才能保证 ≤5 in-flight 仍有序。- 加分区:旧 Key 哈希落点变了,新旧消息可能跨分区被并发消费。
- 加分项:提到 RabbitMQ 是单 Queue 顺序、RocketMQ 是单 MessageQueue 顺序,本质都一样,「最小有序单元」是分区/队列。
Q4:为什么不建议「先开 1 分区,将来需要再加」?
考察点:Kafka 数据模型、运维实践。
答案:
- Key Hash 重分布:加分区后
murmur2(key) % N的结果对几乎所有 Key 都会变,导致同一 Key 的新消息进了新分区。 - 顺序破坏:旧分区里同 Key 的未消费消息 + 新分区里的新消息,被两个 Consumer 并发处理,订单状态机错乱。
- 状态对不上:Kafka Streams / Flink 的 State Store 按 Key 分区存,新消息找不到旧状态。
- 下游幂等失效:用
(topic, partition, offset)当幂等键的下游会出现「同一逻辑消息出现在两个不同 (partition, offset)」。 - 正确做法:
- 一次性按未来 1~2 年容量规划,分区数留 1.5~2 倍余量。
- 真要扩容,重建 Topic(v2、Shadow Topic 双写切换),不是加分区。
- 加分项:提到「Cooperative Sticky Assignor」可以减少 Rebalance 抖动,但不能解决 Key 重分布问题。
Q5:Topic 里 partitions=1000,会出什么问题?
考察点:分区数代价、集群容量。
答案:
- Controller 压力:每个分区的 Leader、ISR、Epoch 都在 Controller 内存,元数据广播延迟变长,故障转移慢。
- 文件句柄爆掉:每分区至少 3 个文件,1000 分区 = 3000 fd,再乘上 Producer/Consumer socket,
ulimit -n不够直接打开失败。 - PageCache 命中率低:1000 个 Segment 的头部都要进 PageCache,冷热混杂,命中率下降。
- Recovery 慢:Broker 重启要重放 1000 分区的 tail 日志,可能 10+ 分钟才上线。
- Producer 端 batch 凑不满:消息分到 1000 分区,每个分区 batch 太小,吞吐反而下降;端到端延迟增大。
- Consumer Rebalance 风暴:分区越多,Rebalance 一次涉及的元数据越多,耗时几十秒,期间消费完全停。
- 正确解法:单 Topic 6~60 个分区为宜;如果业务真的需要 1000 路并行,考虑拆成 10 个 Topic 各 100 分区,或用 RocketMQ / Pulsar 这种分区代价更小的方案。
- 加分项:提到 KRaft 把分区上限拉到了 200 万级别,但元数据传播延迟、Recovery 时间仍然线性增长。
Q6:怎么避免 Kafka 出现「热点分区」?
考察点:Key 设计、负载均衡。
答案:
- 诊断:用 JMX
MessagesInPerSec按 partition 维度看,发现某分区流量是其它的 5~10 倍就是热点。 - 常见原因:
- Key 选成了大客户的
tenantId/ 头部商家的merchantId。 - Key 基数太小(比如只有几十种 Key、却开了 12 分区)。
- Key 为空时早期版本走 Round-Robin 还好,新版 Sticky 会让一段时间内全打到同一分区。
- Key 选成了大客户的
- 修复方案:
- 改 Key 维度:从
tenantId改成orderId,业务允许的话最优。 - 加随机后缀:
key = f"{tenantId}_{random.randint(0, 9)}",把单租户打散到 10 个 Key,但牺牲顺序。 - 大租户隔离:Top N 大租户单独走专用 Topic,小租户走共享 Topic。
- 自定义 Partitioner:在 Producer 端实现
Partitioner接口,让大 Key 主动避开热点分区。
- 改 Key 维度:从
- 数据采样:
kafka-run-class.sh kafka.tools.GetOffsetShell --topic xxx --time -1能看每个分区的 LEO,对比每个分区的写入速度(前后两次差值)。 - 加分项:提到 Pulsar 的「Key_Shared」订阅模式可以在保留 Key 顺序的前提下让一个 Key 被多 Consumer 协作处理,是 Kafka 没有的能力。
本章配套:
06_topic_design/demo.html:分区数 vs 吞吐 / Rebalance 时间 / 元数据量曲线、Key Hash 落点演示、加分区前后落点变化对比。06_topic_design/code/key_distribution.py:模拟 100 万 Key 在不同分区数下的分布均衡度。06_topic_design/code/repartition_impact.py:演示加分区后同一 Key 落点变化的破坏性。
🔗 延伸阅读
- 第 4 章 Producer 深入 —— Partitioner 的源码细节、Sticky Partitioner 的批次策略。
- 第 7 章 存储与日志格式 —— 分区下面是 Segment,分区数变多带来的「文件句柄爆炸」具体长什么样。
- 第 9 章 副本机制与 ISR ——
min.insync.replicas、acks=all、Leader Epoch 的完整故事。 - 第 11 章 消费者组与 Rebalance —— 分区数过多导致的 Rebalance 风暴。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
python
#!/usr/bin/env python3
"""
key_distribution.py
===================
模拟 100 万个 Key 在不同分区数下,经 Murmur2 哈希后的分布均衡度。
输出:
- 每个分区接收的消息数
- 最热分区占比、最冷分区占比
- 变异系数 CV (Coefficient of Variation)
- 分布柱状图(终端 ASCII)
不依赖 Kafka 库,纯 Python 实现 Murmur2 算法(与 Kafka 客户端默认一致)。
用法:
python3 key_distribution.py
python3 key_distribution.py --total 5000000 --partitions 6 12 24 48 --pattern uuid
python3 key_distribution.py --pattern biased # 模拟「大客户占 80% 流量」的场景
python3 key_distribution.py --pattern few # 基数过小(仅 5 个 Key)
"""
from __future__ import annotations
import argparse
import math
import random
import string
import time
import uuid
from collections import Counter
from typing import Iterable, List
# ---------------------------------------------------------------------------
# Kafka 默认 Partitioner 用的 Murmur2 哈希(与 org.apache.kafka.common.utils.Utils#murmur2 等价)
# ---------------------------------------------------------------------------
_M = 0x5BD1E995
_R = 24
_SEED = 0x9747B28C
_MASK32 = 0xFFFFFFFF
def murmur2(data: bytes) -> int:
"""与 Kafka 客户端 Utils.murmur2 完全一致的实现。返回非负 int(toPositive 后)。"""
length = len(data)
h = (_SEED ^ length) & _MASK32
i = 0
while length >= 4:
k = (
(data[i] & 0xFF)
| ((data[i + 1] & 0xFF) << 8)
| ((data[i + 2] & 0xFF) << 16)
| ((data[i + 3] & 0xFF) << 24)
)
k = (k * _M) & _MASK32
k ^= (k >> _R) & _MASK32
k = (k * _M) & _MASK32
h = (h * _M) & _MASK32
h ^= k
i += 4
length -= 4
if length == 3:
h ^= (data[i + 2] & 0xFF) << 16
if length >= 2:
h ^= (data[i + 1] & 0xFF) << 8
if length >= 1:
h ^= data[i] & 0xFF
h = (h * _M) & _MASK32
h ^= (h >> 13) & _MASK32
h = (h * _M) & _MASK32
h ^= (h >> 15) & _MASK32
return h & 0x7FFFFFFF # toPositive
def partition_for(key: str, n: int) -> int:
return murmur2(key.encode("utf-8")) % n
# ---------------------------------------------------------------------------
# Key 生成器
# ---------------------------------------------------------------------------
def gen_keys(pattern: str, total: int) -> Iterable[str]:
"""根据 pattern 生成 Key 流。yield 节省内存。"""
if pattern == "seq":
# user_0 ... user_{total-1}
for i in range(total):
yield f"user_{i}"
elif pattern == "uuid":
for _ in range(total):
yield uuid.uuid4().hex
elif pattern == "tenant":
# 100 个租户,但访问按齐普夫分布——Top 5 占大头
tenants = [f"tenant_{i}" for i in range(100)]
weights = [1.0 / (i + 1) for i in range(100)] # 1, 1/2, 1/3 ... Zipf
s = sum(weights)
weights = [w / s for w in weights]
for _ in range(total):
yield random.choices(tenants, weights=weights, k=1)[0]
elif pattern == "biased":
# 80% 流量打在 5 个大 Key 上,20% 打在普通 Key
hot = [f"VIP_{c}" for c in "ABCDE"]
for i in range(total):
if random.random() < 0.8:
yield hot[i % 5]
else:
yield "normal_" + "".join(
random.choices(string.ascii_lowercase + string.digits, k=10)
)
elif pattern == "few":
for i in range(total):
yield f"group_{i % 5}"
else:
raise ValueError(f"unknown pattern: {pattern}")
# ---------------------------------------------------------------------------
# 统计 + 可视化
# ---------------------------------------------------------------------------
def analyse(counts: List[int]) -> dict:
n = len(counts)
total = sum(counts)
if total == 0:
return {"n": n, "total": 0}
mean = total / n
variance = sum((c - mean) ** 2 for c in counts) / n
stddev = math.sqrt(variance)
cv = stddev / mean if mean else 0
return {
"n": n,
"total": total,
"mean": mean,
"min": min(counts),
"max": max(counts),
"min_pct": min(counts) / total * 100,
"max_pct": max(counts) / total * 100,
"stddev": stddev,
"cv": cv,
"empty": sum(1 for c in counts if c == 0),
}
def print_bar(counts: List[int], width: int = 60) -> None:
"""终端 ASCII 柱状图。横向显示前 N 个分区。"""
if not counts:
return
max_c = max(counts) or 1
avg = sum(counts) / len(counts)
print()
for i, c in enumerate(counts):
bar_len = int(c / max_c * width)
marker = "🔥" if c > avg * 1.5 else " "
bar = "█" * bar_len
print(f" P{i:>3} {marker} {bar} {c:>10,} ({c/sum(counts)*100:5.2f}%)")
print()
# ---------------------------------------------------------------------------
# 主流程
# ---------------------------------------------------------------------------
def run_one(pattern: str, total: int, partitions: int) -> None:
counts = [0] * partitions
t0 = time.time()
for k in gen_keys(pattern, total):
counts[partition_for(k, partitions)] += 1
elapsed = time.time() - t0
stats = analyse(counts)
print(
f"\n========== pattern={pattern} total={total:,} partitions={partitions} =========="
)
print(
f" 耗时 {elapsed:.2f}s | mean={stats['mean']:,.1f}"
f" min={stats['min']:,}({stats['min_pct']:.2f}%)"
f" max={stats['max']:,}({stats['max_pct']:.2f}%)"
)
print(
f" stddev={stats['stddev']:,.1f} CV={stats['cv']:.4f}"
f" empty_partitions={stats['empty']}"
)
if stats["cv"] < 0.05:
verdict = "✅ 极佳:分布几乎完全均匀"
elif stats["cv"] < 0.15:
verdict = "🟢 良好:可放心上线"
elif stats["cv"] < 0.3:
verdict = "🟡 偏斜:可接受,但要监控热点分区"
else:
verdict = "🔴 严重失衡:必须重新设计 Key 或拆 Topic"
print(f" 评估:{verdict}")
print_bar(counts)
def main() -> None:
parser = argparse.ArgumentParser(description="Kafka Key 分布均衡度模拟器")
parser.add_argument("--total", type=int, default=1_000_000,
help="模拟消息总数,默认 100 万")
parser.add_argument("--partitions", type=int, nargs="+",
default=[3, 6, 12, 24, 48],
help="要测试的分区数列表")
parser.add_argument("--pattern", type=str,
choices=["seq", "uuid", "tenant", "biased", "few"],
default="seq",
help="Key 生成模式")
parser.add_argument("--seed", type=int, default=42,
help="随机种子,复现实验")
args = parser.parse_args()
random.seed(args.seed)
print("=" * 70)
print(f" Kafka Key 分布均衡度模拟器")
print(f" pattern={args.pattern} total={args.total:,}")
print(f" partitions={args.partitions}")
print("=" * 70)
for n in args.partitions:
run_one(args.pattern, args.total, n)
print("\n说明:")
print(" CV (Coefficient of Variation) = stddev / mean")
print(" CV ≤ 0.05 极佳")
print(" CV ≤ 0.15 良好")
print(" CV ≤ 0.30 偏斜(需监控)")
print(" CV > 0.30 失衡(建议改 Key)")
if __name__ == "__main__":
main()python
#!/usr/bin/env python3
"""
repartition_impact.py
=====================
演示「加分区」操作对 Key 落点的破坏性。
核心结论:分区数从 N1 → N2 后,
P(同 Key 落点变化) ≈ (N2 - gcd(N1, N2)) / N2
也就是说:分区数加得越多,变化的 Key 就越多。
**这是为什么「加分区」是破坏性操作、必须改用「重建 Topic」的根本原因。**
输出:
- 一份对比表:N1 个分区 vs N2 个分区下,每个 Key 的落点
- 总变化率
- 「按 Key 聚合」的下游业务受影响范围
用法:
python3 repartition_impact.py
python3 repartition_impact.py --n1 6 --n2 12 --keys 1000000
python3 repartition_impact.py --pairs 6:7 6:8 6:12 6:24 12:24 --keys 100000
"""
from __future__ import annotations
import argparse
import math
from collections import Counter, defaultdict
# 复用同目录下的 murmur2
from key_distribution import murmur2, partition_for, gen_keys
def compute_change_rate(n1: int, n2: int, total: int, pattern: str = "seq") -> dict:
"""跑一次实测,返回:
- changed: 落点变化的 Key 数
- same: 落点未变的 Key 数
- rate: 变化率
- moves: {(p1, p2): count} 从 P{p1} 流向 P{p2} 的消息数
"""
moves: dict = defaultdict(int)
changed = same = 0
for k in gen_keys(pattern, total):
p1 = partition_for(k, n1)
p2 = partition_for(k, n2)
moves[(p1, p2)] += 1
if p1 == p2:
same += 1
else:
changed += 1
return {
"n1": n1,
"n2": n2,
"total": total,
"changed": changed,
"same": same,
"rate": changed / total if total else 0,
"moves": moves,
}
def theoretical_rate(n1: int, n2: int) -> float:
"""理论估计:在 Key 哈希均匀的前提下,加分区后的变化比例。
简化模型:n1 → n2 时落点变化概率 ≈ (n2 - gcd(n1, n2)) / n2
(非严格,仅做对比直觉。实际还受 Murmur2 输出分布影响。)
"""
g = math.gcd(n1, n2)
return (n2 - g) / n2
def print_sample(n1: int, n2: int, sample_keys: int = 20) -> None:
"""随机抽 N 个 Key 打印对比表。"""
print(f"\n 样例(前 {sample_keys} 个 user_xxxx):")
print(f" {'Key':<14} {'P (N=' + str(n1) + ')':<10} {'P (N=' + str(n2) + ')':<10} {'变化?':<6}")
print(f" {'-'*14} {'-'*10} {'-'*10} {'-'*6}")
for i in range(sample_keys):
k = f"user_{i}"
p1 = partition_for(k, n1)
p2 = partition_for(k, n2)
flag = "🔴 变" if p1 != p2 else "🟢 同"
print(f" {k:<14} {p1:<10} {p2:<10} {flag}")
def print_top_moves(result: dict, top: int = 8) -> None:
"""打印「最常发生的迁移路径 P_i → P_j」"""
moves = sorted(result["moves"].items(), key=lambda kv: -kv[1])[:top]
total = result["total"]
print(f"\n 迁移热力图(最常见 Top-{top}):")
print(f" {'From':<8} {'To':<8} {'Count':>10} {'Pct':>8}")
print(f" {'-'*8} {'-'*8} {'-'*10} {'-'*8}")
for (p1, p2), c in moves:
flag = " " if p1 == p2 else "→ "
print(f" P{p1:<6} {flag}P{p2:<6} {c:>10,} {c/total*100:>7.2f}%")
def downstream_impact(result: dict) -> None:
"""评估「按 Key 聚合」的下游业务受影响范围。"""
rate = result["rate"]
n1, n2 = result["n1"], result["n2"]
print()
print(f" ┌── 下游受影响评估 ──────────────────────────────────────")
print(f" │ 实测变化率: {rate*100:.2f}%")
print(f" │ 理论估计 : {theoretical_rate(n1, n2)*100:.2f}% (公式: (n2-gcd)/n2)")
if rate > 0.5:
print(f" │ 等级 : 🔴 重度破坏(>50% 的 Key 迁移)")
print(f" │ 影响范围 : Kafka Streams / Flink 状态、Redis 按 Key 缓存、")
print(f" │ 所有按 Key 聚合的下游统计 —— 全部失效。")
print(f" │ 建议 : 不要加分区。建新 Topic learn.06.xx.v2,双写切流。")
elif rate > 0.2:
print(f" │ 等级 : 🟡 部分破坏(20%~50% 的 Key 迁移)")
print(f" │ 影响范围 : 同上,但比例较小。仍需通知所有下游做兼容评估。")
print(f" │ 建议 : 评估业务能否容忍统计误差,否则同样建议重建 Topic。")
else:
print(f" │ 等级 : 🟢 较小(<20%)")
print(f" │ 注意 : 即使比例小,仍存在「同 Key 在新旧分区上短期并存」")
print(f" │ 的风险,会导致瞬时顺序错乱。")
print(f" └─────────────────────────────────────────────────────")
def run_pair(n1: int, n2: int, total: int, pattern: str) -> None:
print(f"\n{'='*72}")
print(f" 分区数变更:{n1} → {n2} (总消息数 {total:,}, pattern={pattern})")
print(f"{'='*72}")
result = compute_change_rate(n1, n2, total, pattern)
print_sample(n1, n2, sample_keys=20)
print_top_moves(result, top=10)
downstream_impact(result)
def main() -> None:
parser = argparse.ArgumentParser(description="Kafka 加分区破坏性演示")
parser.add_argument("--n1", type=int, default=6, help="分区数(前)")
parser.add_argument("--n2", type=int, default=8, help="分区数(后)")
parser.add_argument("--keys", type=int, default=100_000, help="模拟消息数")
parser.add_argument("--pairs", type=str, nargs="*", default=None,
help="多组对比,格式 n1:n2,如 6:7 6:8 6:12 12:24")
parser.add_argument("--pattern", type=str, default="seq",
choices=["seq", "uuid", "tenant", "biased", "few"])
args = parser.parse_args()
if args.pairs:
for p in args.pairs:
n1, n2 = (int(x) for x in p.split(":"))
run_pair(n1, n2, args.keys, args.pattern)
else:
run_pair(args.n1, args.n2, args.keys, args.pattern)
print("\n小结:")
print(" - Kafka 默认 Partitioner = murmur2(key) % N")
print(" - 分区数变化时,几乎所有 Key 的落点都会变。")
print(" - 唯一例外是「N1 是 N2 的因子」时,部分 Key 仍会落到 P_i 或 P_(i+N1),")
print(" 但仍有 (n2-gcd)/n2 的比例发生迁移。")
print(" - 因此:生产环境想扩容请「重建 Topic」,不要直接 alter --partitions。")
if __name__ == "__main__":
main()