Skip to content

第 6 章 Topic 设计与分区策略:从「能跑」到「跑得稳」

目标读者:会用 kafka-topics.sh --create 拍脑袋建 Topic、副本数永远填 1、分区数永远填 1 或者上百个,然后线上要么扛不住要么 Rebalance 抖死的同学。

学完你会:拿到一个新业务,能在 5 分钟内说清楚「这个 Topic 应该开多少分区、多少副本、Key 怎么定、min.insync.replicas 填几、未来扩容怎么办」,并且讲得出每一项背后的依据。


0. 导读:分区数和副本数,是 Kafka 工程师的「灵魂两问」

绝大多数线上 Kafka 事故的根因,可以归结成下面这几条之一:

  1. 分区数选少了 —— 后端再多消费者也吃不动,records-lag-max 一路涨;想加分区,但 Key Hash 分布会变,下游所有「按 Key 聚合」的逻辑全乱。
  2. 分区数选多了 —— 单 Broker 上几万个 Partition,Controller 元数据爆炸、OpenFileDescriptors 顶到上限、宕机后 Recovery 几十分钟,新 Leader 怎么也选不出来。
  3. 副本数填了 1 —— 一台 Broker 宕机,整个 Topic 直接 OFFLINE,业务方电话打爆。
  4. acks=1 + min.insync.replicas=1 —— 看着「保留 3 副本」,其实写下去只要 Leader 自己确认就 ack,Leader 一挂数据丢了一片。
  5. Key 用了 userId 但热门用户占 80% 流量 —— 单个分区直接被打爆,其它分区闲着发呆。
  6. 「先小心试运行」副本数 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/skafka-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=allcompression=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:9092

Kafka 会自动让 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")   # 没传 key

Kafka 2.4+ 默认走 Sticky Partitioner

  • 选定一个分区,整个 batch 都发同一个分区
  • 等 batch 被发出去(达到 linger.msbatch.size)之后,再换一个分区。

老版本是「Round-Robin」每条换分区,会让每个分区只攒到一条消息,batch 永远凑不满。Sticky 之后吞吐能提升 2 ~ 3 倍。

没 Key 的副作用:消息分布是「按时间分块」的,不是「按 Key 哈希分均匀」的。如果业务真的要均匀分布且不需要顺序,可以加一个随机 Key

python
producer.send("topic", key=str(uuid.uuid4()).encode(), value=value)

5. min.insync.replicasacks 的配合

5.1 acks 的三档语义

acksProducer 何时收到 ack一致性性能
0发出去就算成功,不等任何 ack可能丢消息(Broker 没收到也不知道)最快
1Leader 写入本地日志(不一定 fsync)就 ackLeader 突然挂掉、还没复制给 Follower → 丢消息中等
all(= -1ISR 中所有副本都已经写入本地日志才 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 三句话讲清楚

  1. 单分区内:消息 100% 严格按写入顺序,Consumer poll 出来也 100% 按顺序。
  2. 跨分区:Kafka 不提供任何顺序保证。
  3. 跨 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%。具体来说:

Keymurmur2 % 6murmur2 % 8落点是否变化
order_001 (hash=12345678)12345678 % 6 = 412345678 % 8 = 6变了
order_002 (hash=98765432)98765432 % 6 = 298765432 % 8 = 0变了
order_003 (hash=11111111)11111111 % 6 = 511111111 % 8 = 7变了

大部分 Key 都会换分区

7.2 谁会受伤?

  1. 消费者侧的「按 Key 聚合」状态
    • 例:Kafka Streams 的 aggregate(key)order_001 的状态存在 partition 4 对应的 State Store 里。加分区后 order_001 的新消息去了 partition 6,但旧状态还在 partition 4 的 Store,状态对不上号。
  2. 顺序保证被打破
    • 加分区瞬间,order_001 的新消息进了 partition 6,但 partition 4 上还有一条没消费完的旧消息。Consumer A(消费 P4)和 Consumer B(消费 P6)并发处理,旧消息可能在新消息之后才被处理。
  3. 下游去重 / 幂等表的依赖
    • 如果下游用 (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.v1
  • prod.order.purchase.paid.v2
  • dev.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

逐项解读:

配置说明
partitions12业务峰值 200MB/s ÷ 单分区 40MB/s ≈ 5,留 2 倍余量
replication.factor3黄金三角
min.insync.replicas2黄金三角
retention.ms6048000007 天
segment.bytes512MB控制单 Segment 大小,避免太多小文件
compression.typeproducer透传 Producer 端压缩,不做解压再压
max.message.bytes1MB防止巨型消息撑爆 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 的对比

维度KafkaRabbitMQRocketMQPulsar
分区概念Topic-Partition 显式Queue 不分区,靠多 Queue + ExchangeMessageQueue(类似 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 哈希到分区
推荐 RF33 节点镜像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:分区数怎么选?给一个公式和一个例子。

考察点:吞吐建模、容量规划。

答案

  1. 公式partitions = max(T/P_p, T/C_c) × 1.5~2 余量,其中 T 是业务目标吞吐,P_p 是单分区生产吞吐(实测),C_c 是单消费者吞吐(实测)。
  2. :目标 100MB/s,单分区写 40MB/s,单消费者 20MB/s → max(100/40, 100/20) = max(3, 5) = 5,再 ×2 余量 = 10 个分区
  3. 上限考量:单 Broker 一般不超 4000 分区,否则 Controller 元数据、文件句柄、Recovery 时间都会炸。
  4. 不能无脑大:分区越多,每个 batch 越凑不满,端到端延迟反而上升;同时加分区会让 Key Hash 重分布,破坏顺序。
  5. 加分项:提到 Sticky Partitioner(无 Key 时)让 batch 集中、不至于分得太散;提到 Kafka 4.x KRaft 下分区上限提高到几十万但仍要看具体硬件。

Q2:副本数 3 + min.insync.replicas=2 + acks=all 这套组合,能保证什么、不能保证什么?

考察点:副本机制、可靠性语义。

答案

  1. 能保证:消息一旦 ack,至少在 2 个 Broker 的本地日志里。容忍 1 台 Broker 同时宕机不丢数据。
  2. 不能保证:ISR 缩到 1 个时 Producer 拒写(NotEnoughReplicasException),所以业务方要做好重试和报警。
  3. acks=all 的细节:是「ISR 中所有副本」ack,不是「所有副本」。所以必须配 min.insync.replicas 才安全。
  4. 故障行为
    • 1 台 Broker 挂:写入仍然成功(ISR 还剩 2)。
    • 2 台 Broker 挂(且 Leader 是其中之一):触发 Leader 选举,几秒不可用,但不丢数据。
    • 全集群崩:完全不可用,但只要重启起来数据还在。
  5. 加分项:提到 unclean.leader.election.enable=false(默认 false)保证不会让 Out-of-Sync 副本当 Leader 导致 ack 过的消息被截断;提到 enable.idempotence=true 解决 retries 引起的重复。

Q3:什么情况下消息能保证顺序?跨分区怎么办?

考察点:顺序语义、Key 设计。

答案

  1. 单分区顺序:同一个分区的消息,写入顺序 == 存储顺序 == 消费顺序。
  2. 跨分区无序:Kafka 不提供任何跨分区顺序保证。
  3. 怎么落到同分区:Producer 设 Key,Kafka 用 murmur2(key) % num_partitions 决定分区。Key 相同 → 分区相同 → 顺序保证。
  4. 业务层面的解决方案
    • 同一聚合根(订单号、用户 id)用同一 Key,让相关事件同分区。
    • 必要时按版本号 / 时间戳在消费者端排序。
    • 极端情况下用单分区 Topic 实现全局有序(牺牲吞吐)。
  5. 顺序的隐藏破坏者
    • max.in.flight.requests.per.connection > 1 + 关掉幂等:retry 会乱序。开了幂等才能保证 ≤5 in-flight 仍有序。
    • 加分区:旧 Key 哈希落点变了,新旧消息可能跨分区被并发消费。
  6. 加分项:提到 RabbitMQ 是单 Queue 顺序、RocketMQ 是单 MessageQueue 顺序,本质都一样,「最小有序单元」是分区/队列。

Q4:为什么不建议「先开 1 分区,将来需要再加」?

考察点:Kafka 数据模型、运维实践。

答案

  1. Key Hash 重分布:加分区后 murmur2(key) % N 的结果对几乎所有 Key 都会变,导致同一 Key 的新消息进了新分区。
  2. 顺序破坏:旧分区里同 Key 的未消费消息 + 新分区里的新消息,被两个 Consumer 并发处理,订单状态机错乱。
  3. 状态对不上:Kafka Streams / Flink 的 State Store 按 Key 分区存,新消息找不到旧状态。
  4. 下游幂等失效:用 (topic, partition, offset) 当幂等键的下游会出现「同一逻辑消息出现在两个不同 (partition, offset)」。
  5. 正确做法
    • 一次性按未来 1~2 年容量规划,分区数留 1.5~2 倍余量。
    • 真要扩容,重建 Topic(v2、Shadow Topic 双写切换),不是加分区。
  6. 加分项:提到「Cooperative Sticky Assignor」可以减少 Rebalance 抖动,但不能解决 Key 重分布问题。

Q5:Topic 里 partitions=1000,会出什么问题?

考察点:分区数代价、集群容量。

答案

  1. Controller 压力:每个分区的 Leader、ISR、Epoch 都在 Controller 内存,元数据广播延迟变长,故障转移慢。
  2. 文件句柄爆掉:每分区至少 3 个文件,1000 分区 = 3000 fd,再乘上 Producer/Consumer socket,ulimit -n 不够直接打开失败。
  3. PageCache 命中率低:1000 个 Segment 的头部都要进 PageCache,冷热混杂,命中率下降。
  4. Recovery 慢:Broker 重启要重放 1000 分区的 tail 日志,可能 10+ 分钟才上线。
  5. Producer 端 batch 凑不满:消息分到 1000 分区,每个分区 batch 太小,吞吐反而下降;端到端延迟增大。
  6. Consumer Rebalance 风暴:分区越多,Rebalance 一次涉及的元数据越多,耗时几十秒,期间消费完全停。
  7. 正确解法:单 Topic 6~60 个分区为宜;如果业务真的需要 1000 路并行,考虑拆成 10 个 Topic 各 100 分区,或用 RocketMQ / Pulsar 这种分区代价更小的方案。
  8. 加分项:提到 KRaft 把分区上限拉到了 200 万级别,但元数据传播延迟、Recovery 时间仍然线性增长。

Q6:怎么避免 Kafka 出现「热点分区」?

考察点:Key 设计、负载均衡。

答案

  1. 诊断:用 JMX MessagesInPerSec 按 partition 维度看,发现某分区流量是其它的 5~10 倍就是热点。
  2. 常见原因
    • Key 选成了大客户的 tenantId / 头部商家的 merchantId
    • Key 基数太小(比如只有几十种 Key、却开了 12 分区)。
    • Key 为空时早期版本走 Round-Robin 还好,新版 Sticky 会让一段时间内全打到同一分区。
  3. 修复方案
    • 改 Key 维度:从 tenantId 改成 orderId,业务允许的话最优。
    • 加随机后缀key = f"{tenantId}_{random.randint(0, 9)}",把单租户打散到 10 个 Key,但牺牲顺序
    • 大租户隔离:Top N 大租户单独走专用 Topic,小租户走共享 Topic。
    • 自定义 Partitioner:在 Producer 端实现 Partitioner 接口,让大 Key 主动避开热点分区。
  4. 数据采样kafka-run-class.sh kafka.tools.GetOffsetShell --topic xxx --time -1 能看每个分区的 LEO,对比每个分区的写入速度(前后两次差值)。
  5. 加分项:提到 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.replicasacks=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()

key_distribution.py ↗ · repartition_impact.py ↗