Producer 分区策略
定义与作用
分区策略是 Producer 决定每条消息发往哪个 Partition 的规则。它的核心价值在于:让应用开发者通过 Key 控制消息的路由,从而实现业务级别的有序性保证——如同一用户的消息进入同一分区、同一订单的状态变更有序处理。
核心原理
分区器决策流程
三种分区策略对比
| 策略 | 触发条件 | 分布特点 | 有序性 |
|---|---|---|---|
| 显式指定 | 代码中设置 partition | 由开发者控制 | 完全由开发者决定 |
| Key 哈希 | 设置 key | 相同 Key → 同一分区 | Key 级别有序 |
| Sticky(2.4+ 默认) | 无 Key | 批次内粘滞、批次间切换 | 无序 |
| Round-Robin(旧版默认) | 无 Key | 每条消息轮询切换 | 无序 |
为什么 Round-Robin 被 Sticky 取代?
Round-Robin:每条消息换分区
msg1→P0, msg2→P1, msg3→P2, msg4→P0, msg5→P1, ...
问题:每条消息单独一个请求,批处理完全失效。
Sticky Partitioner:攒一批再换
batch1: [msg1, msg2, ..., msg100] → P0
batch2: [msg101, msg102, ..., msg200] → P1
...
优势:批处理生效,网络效率高。
完整示例
示例一:Key 哈希实现用户级别消息有序
场景:用户行为日志系统。每个用户的行为(浏览、加购、下单、支付)需要严格按时间顺序处理。
操作前环境:Topic user-behavior,3 分区。
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
// 用户 1001 的行为序列
producer.send(new ProducerRecord<>("user-behavior", "1001", "view_item:book"));
producer.send(new ProducerRecord<>("user-behavior", "1001", "add_to_cart:book"));
producer.send(new ProducerRecord<>("user-behavior", "1001", "checkout:order-5001"));
// 用户 2002 的行为序列
producer.send(new ProducerRecord<>("user-behavior", "2002", "view_item:phone"));
producer.send(new ProducerRecord<>("user-behavior", "2002", "add_to_cart:phone"));
producer.flush();
}
操作后验证:
# 查看各分区的最新 Offset
bin/kafka-run-class.sh kafka.tools.GetOffsetShell \
--topic user-behavior --time -1 --bootstrap-server localhost:9092
# user-behavior:0:3 ← 用户 1001 的 3 条消息
# user-behavior:1:2 ← 用户 2002 的 2 条消息
# user-behavior:2:0
操作前后对比:
| 用户 ID | 消息数 | 所在分区 | 分区内顺序 |
|---|---|---|---|
| 1001 | 3 | Partition 0 | view → add_to_cart → checkout(严格有序) |
| 2002 | 2 | Partition 1 | view → add_to_cart(严格有序) |
示例二:自定义 Partitioner
场景:VIP 用户(ID 以 V 开头)的消息路由到特定分区,便于优先处理。
import org.apache.kafka.clients.producer.Partitioner;
import org.apache.kafka.common.Cluster;
public class VipPartitioner implements Partitioner {
@Override
public int partition(String topic, Object key, byte[] keyBytes,
Object value, byte[] valueBytes, Cluster cluster) {
int partitionCount = cluster.partitionCountForTopic(topic);
String keyStr = (String) key;
if (keyStr != null && keyStr.startsWith("V")) {
// VIP 用户 → 固定到最后一个分区
return partitionCount - 1;
}
// 普通用户 → 默认哈希
return Math.abs(keyStr.hashCode()) % (partitionCount - 1);
}
@Override
public void close() {}
@Override
public void configure(Map<String, ?> configs) {}
}
配置使用:
props.put("partitioner.class", "com.example.VipPartitioner");
// 发送消息
producer.send(new ProducerRecord<>("orders", "V1001", "VIP order")); // → Partition 2
producer.send(new ProducerRecord<>("orders", "U2002", "normal order")); // → Partition 0 或 1
操作前后对比:
| Key | 默认哈希分区 | 自定义 VIP Partitioner |
|---|---|---|
V1001 | 随机(0-2) | 固定 Partition 2 |
U2002 | 随机(0-2) | 在 0-1 之间哈希 |
易错场景
易错 1:Key 的哈希分布不均衡
场景:使用简单的 String.hashCode() 且 Key 集中在少数几个值。
后果:热点分区——某个分区数据量远大于其他分区,导致该 Broker 磁盘、CPU 过载。
缓解方法:
- Kafka 内置使用 Murmur2 哈希(分布比 Java
hashCode均匀) - 增加分区数(分散到更多分区)
- 考虑在 Key 中加入时间戳或序列号提高分散度
易错 2:分区数变更后 Key 路由被打乱
场景:Topic 从 3 分区增加到 6 分区。
后果:murmur2(key) % 3 和 murmur2(key) % 6 结果不同。原本在 Partition 0 的某 Key 可能被路由到 Partition 3。
影响:该 Key 的新消息进入新分区,而消费者可能还在旧分区等待——造成该 Key 的消息在逻辑上"乱序"。
最佳实践:分区数的规划一开始就要做足,尽量避免后期频繁增加。
面试高频考点
Q:Kafka 默认的分区器在选择分区时的优先级是什么?
A:优先级从高到低:
- 显式指定 partition — 代码中直接指定 > 一切
- Key 的 Murmur2 哈希 —
murmur2(key_bytes) % partition_count - Sticky Partitioner(无 Key 时) — 批次粘滞,批次满后切换
partition 和 key 同时指定时,以 partition 为准。