Learn
Kafka/10-partitioning-ordering

分区策略与消息顺序

"消息发到哪个分区"看似小事,实际决定了三件大事:顺序保证、负载均衡、消费并行度。本章讲清分区路由规则、如何保证业务需要的顺序,以及"分区数到底设多少"这个永恒问题。

1. 消息路由规则

Producer 决定分区的优先级:

1. ProducerRecord 显式指定了 partition  -> 用它
2. 有 key                              -> murmur2(key) % 分区数
3. 无 key                              -> 粘性分区 (sticky partitioning)

1.1 按 key 哈希

// key = "user-42" 的所有消息永远进同一个分区
producer.send(new ProducerRecord<>("orders", "user-42", payload));

同 key 必然同分区,而分区内严格有序——这就是 Kafka 顺序保证的基石。

1.2 无 key:粘性分区

2.4 之前无 key 消息逐条轮询分区,导致每个批次都很小(一批只归属一个分区)。粘性分区改为:

轮询 (旧):  m1->p0  m2->p1  m3->p2  m4->p0 ...   批次碎, 吞吐差
粘性 (新):  m1..m50 都 -> p0 (直到批次满/linger到期)
           下一批 -> p2 (随机换一个分区)
           长期看各分区仍均匀, 但每批都是满批

效果:无 key 场景延迟与吞吐显著改善,长期分布依然均匀。

2. 自定义 Partitioner

需要特殊路由时(如大客户隔离到专属分区)实现 Partitioner 接口:

public class VipPartitioner implements Partitioner {
    @Override
    public int partition(String topic, Object key, byte[] keyBytes,
                         Object value, byte[] valueBytes, Cluster cluster) {
        int numPartitions = cluster.partitionCountForTopic(topic);
        String k = (String) key;
        if (k != null && k.startsWith("vip-")) {
            return numPartitions - 1;              // VIP 独占最后一个分区
        }
        // 其余 key 哈希到前 n-1 个分区
        return Math.abs(Utils.murmur2(keyBytes)) % (numPartitions - 1);
    }
    @Override public void close() {}
    @Override public void configure(Map<String, ?> configs) {}
}
props.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, VipPartitioner.class.getName());

Python(confluent-kafka)不支持自定义分区器类,但可以在调用 produce 时显式指定分区实现同样效果:

def pick_partition(key: str, num_partitions: int) -> int:
    if key.startswith("vip-"):
        return num_partitions - 1
    return hash(key) % (num_partitions - 1)
 
producer.produce("orders", key=key, value=payload,
                 partition=pick_partition(key, 6))

3. 顺序保证:能保什么、不能保什么

Kafka 的顺序承诺只有一条:单分区内,消息按写入顺序存储和消费。

需求可行方案
全局有序单分区 topic(吞吐受限于单分区,一般只用于低流量场景)
同一业务实体有序(同一订单、同一用户)以实体 ID 为 key → 局部有序,绝大多数业务的正确答案
完全不关心顺序无 key,让粘性分区做负载均衡

局部有序还需要 Producer 端配合,否则重试会打乱同分区内的顺序:

enable.idempotence=true    # 3.x 默认开启, 同时解决重试乱序与重复
# 幂等开启时 max.in.flight <= 5 都能保证分区内有序

消费端的隐形陷阱:消费者拿到消息后如果丢进多线程池处理,顺序就没了。要保顺序,同一分区的消息必须串行处理(或按 key 再次分发到固定线程)。

端到端顺序 = 同 key 同分区 (Producer)
           + 幂等防重试乱序 (Producer)
           + 分区内串行处理 (Consumer)
三环缺一不可

4. 分区数怎么选

分区数是并行度的上限,也是元数据与文件句柄的开销来源。

4.1 估算公式

分区数 = max( 目标吞吐 / 单分区生产吞吐,
              目标吞吐 / 单分区消费吞吐,
              预期最大消费者数 )

经验值:单分区生产吞吐 10-50 MB/s(取决于压缩与磁盘),单分区消费吞吐通常受业务处理速度限制。

4.2 实用建议

  • 中小流量业务 topic:6 或 12 起步(能被 2/3/4/6 整除,扩消费者灵活)。
  • 预留增长空间:按 1-2 年后的峰值估算,因为扩分区有副作用(见下)。
  • 单 broker 承载分区数(含副本)建议几千以内;集群总分区数过多会拖慢选举与恢复。
  • 不要为"以后可能的大流量"直接开几百个分区:小流量 + 超多分区会导致批次碎、端到端延迟升高。

4.3 扩分区的影响

分区只能增不能减。扩分区的代价:

  1. key 路由改变:hash(key) % N 中 N 变了,同 key 的新消息落到新分区,历史顺序语义被打破。
  2. 已有数据不迁移:老消息留在原分区,新消息按新规则分布,短期数据倾斜。
  3. 触发消费组 Rebalance。

依赖 key 顺序的业务扩分区标准流程:暂停生产 → 等消费清空积压 → 扩分区 → 恢复。或者一步到位建新 topic 双写迁移。

⚠️分区倾斜: 热 key 问题

按 key 哈希均匀的前提是 key 本身分布均匀。如果 20% 的流量来自同一个大客户 ID,它所在的分区就会成为热点——生产端批次堆积、消费端 LAG 单分区飙高。解法:换更细粒度的 key(如 订单ID 而非 客户ID,如果顺序要求允许)、大客户单独 topic、或自定义分区器把热 key 打散(牺牲该 key 的顺序)。

5. 验证工具

# 看每个分区的消息量是否均匀 (LEO - earliest)
bin/kafka-get-offsets.sh --bootstrap-server localhost:9092 --topic orders --time -1
bin/kafka-get-offsets.sh --bootstrap-server localhost:9092 --topic orders --time -2
 
# 消费时打印分区验证路由
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic orders --from-beginning \
  --property print.key=true --property print.partition=true
💡key 的设计原则

选 key 时问自己两个问题:1)哪些消息之间必须有序?它们应该共享同一个 key。2)这个 key 的取值分布够分散吗?满足顺序要求的前提下,key 粒度越细、分布越均匀越好。"userId" 通常是顺序与均匀的最佳平衡点。

小结

  • 路由优先级:显式分区 > key 哈希 > 粘性分区;同 key 必同分区
  • 粘性分区解决无 key 消息批次碎的问题,长期分布依然均匀
  • 顺序三环:同 key 同分区 + Producer 幂等 + 消费端分区内串行
  • 分区数按吞吐和消费者数估算,6/12 起步、留增长空间;只能增不能减
  • 扩分区会改变 key 路由、破坏历史顺序语义,热 key 会造成分区倾斜
🎯练习
  1. 向 4 分区 topic 分别发送 1000 条有 key(10 个 key)和 1000 条无 key 消息,用 get-offsets 统计各分区消息量,观察两种模式的分布。
  2. 实现一个自定义 Partitioner,把以 "audit-" 开头的 key 全部路由到 0 号分区,其余正常哈希,并写消费者验证。
  3. 把 3 分区扩到 6 分区,继续用同一批 key 生产,验证同一 key 扩容前后落在不同分区。