前一篇立了地基:Kafka 是一份分布式可重放日志。这一篇聚焦在 Producer 端一个看起来微不足道、实则牵一发动全身的决策:
一条消息,到底该写进哪个 partition?
这个决策(分区策略)同时决定了吞吐、延迟、顺序性。而 Kafka 2.4 引入的 Sticky Partitioner 就是一次针对这个决策的精妙优化——在几乎不改变数据分布的前提下,把无 key 消息的 p99 延迟砍掉约一半。下面把来龙去脉讲透。(分区写进去后”分给哪个消费者”的另一半——消费侧 assignor——放在消费者篇 §4 讲。)
TL;DR
- 发送不是”一条一发”,而是”攒批再发”:消息先进
RecordAccumulator按分区攒成批次,batch.size写满或linger.ms到点才由 Sender 发出。批越大、批数越少,单条成本越低、延迟越低。 - 分区策略决定了”能不能攒成一批”:只有落到同一分区的消息才能进同一个 batch。分区选得不好,批就攒不起来。
- 四种生产者策略:默认(有 key 按 key 哈希 / 无 key 走分区器)、轮询(round-robin)、Sticky(黏性)、自定义。
- 老的无 key 策略(2.4 前)= 逐条轮询:把无 key 消息一条条撒到不同分区,导致大量小批次 → 更多请求 → 更高延迟。
- Sticky Partitioner(2.4+):无 key 消息先粘住一个分区猛灌,直到该分区的批次写满/完成,再随机换下一个分区。长期看数据依旧均匀分布,但批更大、批数更少。
- Confluent 基准:16 分区、3 producer、1000 msg/s 场景下,sticky 的 p99 延迟约为默认策略的一半;分区越多、
linger.ms越大,优势越明显,还常降低 CPU。 - 对有 key 消息无影响:有 key 仍按 key 哈希路由,sticky 只优化无 key 的批处理;实测有 key / 混合 key 场景零副作用。
- key→分区映射不是永久稳定的:哈希只在分区数不变时稳定;一旦扩分区,同一 key 可能落到新分区,历史顺序保证被打破。
- 顺序性铁律不变:Kafka 只保证分区内有序;要顺序就用 key 把相关消息钉到同一分区。要严格幂等有序还需配
enable.idempotence等参数。 - 写进去只是一半:分区写进去后”分给哪个消费者”由消费端的 assignor(Range / RoundRobin / Sticky)决定,那属于消费链路——拆解见消费者篇 §4。
Table of contents
Open Table of contents
1. 先看生产者发送流程:一切从”攒批”说起
要理解分区策略为什么重要,得先知道 Producer 不是收到一条就发一条。先把要优化的指标定义清楚:
生产者延迟 ≈ 一条消息从被客户端
send()到**被 Kafka 确认(acknowledged)**所经历的时间。降低这个时间,整条链路都受益——而它的大头,往往花在”攒批与请求排队”上。
一条消息的旅程:
- 封装成
ProducerRecord→ 序列化; - 交给 Partitioner 算出目标分区;
- 写进
RecordAccumulator,按分区攒成RecordBatch; - Sender 线程把攒好的批次封成
Request批量发往 Broker。
批次的完成条件只有两个:
batch.size(默认 16384 字节)写满 → 立刻发;或linger.ms(默认 0)时间到 → 发。 谁先到听谁的。
这里有个反直觉但关键的事实:即使 linger.ms=0,也不代表每条消息单独成批。 因为系统处理每个请求本身需要一点点时间,在这段时间里同时到达、且发往同一分区的消息会被自然地攒进同一个 batch。换句话说,攒批不完全靠等待时间,也靠”处理请求的固有间隙”。
而”能不能攒进同一个 batch”的前提,是这些消息去了同一个分区——这就把球踢给了分区策略。每处理一个 batch 都有固定开销,批里记录越少,单条记录分摊的成本越高;小批次多 → 请求多、排队多 → 延迟高。记住这条因果链,后面所有优化都是围着它转。
2. 四种分区策略
Kafka 允许你通过配置 partitioner.class 来决定每条记录进哪个分区。开箱四种:
| 策略 | 行为 | 适用场景 |
|---|---|---|
| 默认分区器(Default) | 有 key:对 key 做哈希映射到分区(同 key → 同分区)。无 key:交给分区器分配(2.4 前逐条轮询,2.4 后走 sticky)。 | 绝大多数场景的缺省选择 |
| 轮询分区器(Round-robin) | 消息以轮询方式依次撒到各分区,无视 key,力求绝对均匀。 | key 高度倾斜、又不需要按 key 保序时 |
| 黏性分区器(Uniform Sticky) | 无 key 消息粘住一个分区发,直到 batch.size 满或 linger.ms 到,再换分区——为降低延迟而生。 | 无 key 高吞吐(2.4+ 已内置进默认器) |
| 自定义分区器(Custom) | 实现 Partitioner 接口,在 partition() 里写自己的 key→分区路由逻辑。 | 热点 key 隔离、按业务维度定向 |
真正的分歧点在无 key 消息怎么办。有 key 的情况所有策略基本一致(按 key 哈希),没有争议;问题全出在无 key 消息的分配方式上。
不过有 key 的路由也藏着两个容易被忽略的细节:
2.1 key 哈希的隐藏前提:分区数不能变
同一个 key 落到同一个分区,靠的是 hash(key) % 分区数。这意味着这条稳定性只在分区数不变时成立:
一旦给 topic 扩容分区,取模的除数变了,同一个 key 可能被路由到与历史不同的新分区——从此该 key 的新老消息分散在两个分区,“同 key 分区内有序”的保证就被打破了。
所以”按 key 保序”的系统要把分区数当成近乎不可变的契约来对待:要么一开始就按目标吞吐把分区数留足(见核心原理篇的分区容量决策),要么在扩分区时接受一段顺序性的”缝合期”。
2.2 轮询不是没用:它是”抗 key 倾斜”的解药
轮询常被当成”过时策略”,但它有个默认器给不了的价值:当负载被单个热点 key 压垮时。如果用默认哈希,某个超热 key(比如占 40% 流量的大客户)的消息会全部涌向同一个分区,造成该分区堆积、对应消费者过载、broker 空间倾斜。此时若不要求按 key 保序,轮询能把写入均摊到所有分区,反而更健康。
换句话说:sticky 优化的是”无 key 的批处理效率”,轮询解决的是”有 key 但严重倾斜的负载均衡”——两者针对的问题不同,别混为一谈。
3. 问题:老的”逐条轮询”为什么慢
2.4 之前,默认分区器处理无 key 消息的方式是逐条轮询——来一条,发一个分区;再来一条,发下一个分区……看起来最公平,但它和”攒批”是天生的冤家:
问题就在这张图的上半部分:逐条轮询把消息摊得到处都是,每个分区同一时刻只攒到寥寥几条,于是形成一大堆小批次。 而:
小批次 = 更多请求 + 更多排队 = 更高延迟。 每处理一个 batch 都有固定开销,批里记录越少,单条记录分摊的成本越高。
于是”为了均匀而逐条轮询”这个朴素直觉,反而在批处理层面帮了倒忙——尤其当分区数很多、吞吐不高时,消息还没来得及在某个分区攒够一批,就被轮询规则赶到下一个分区去了。
4. Sticky Partitioner:粘住一个分区,把批攒大
Kafka 2.4(KIP-480)给出的解法优雅得近乎”作弊”:
黏性分区器把所有无 key 消息都发往同一个”被选中的粘性分区”,直到该分区的批次被填满或完成,才随机挑选并”粘”到一个新分区。 这样在更长的时间尺度上,记录依旧大致均匀分布到所有分区,同时享受到更大批次的红利。
实现上,2.4 给 partitioner 接口加了个 onNewBatch 方法,在新批次即将创建前被调用——正是切换粘性分区的最佳时机,DefaultPartitioner 实现了它。
核心洞察就一句话:把”横向撒开”改成”纵向灌满再换”。 短期集中火力把一个分区的批攒大,长期通过不断换分区保持整体均匀——鱼和熊掌兼得。
对照图的下半部分:同样 6 条消息,sticky 先把 P0 灌到批满(4 条)再换 P1,批次数从 3 个降到更少、每个批更大,请求数随之下降。
5. Confluent 的基准数据:到底快多少
光讲原理不够。Confluent 用 Kafka 自带的压测框架 Trogdor(配合 Castle 测试编排)在 AWS 上跑了一组 ProduceBench,把方法学摊开来看结论才站得住脚。
5.1 测试设定(可复现)
| 项 | 取值 |
|---|---|
| 测试时长 | 12 分钟 |
| Broker 数 | 3(AWS m3.xlarge + SSD) |
| Producer 数 | 1–3 |
replication factor | 3 |
| Topic 数 | 4 |
linger.ms | 0(除单独说明的 linger 实验外) |
acks | all |
| key 生成 | 全部为 null(考察无 key 路径) |
useConfiguredPartitioner / skipFlush | true / true |
两个
true是关键:useConfiguredPartitioner确保用DefaultPartitioner真正决定分区,skipFlush确保批次是靠写满或linger.ms触发而不是被 flush 逼出——这样对比的才是”分区策略本身”的效果。
5.2 核心结论
- p99 延迟砍半:3 producer、每秒 1000 条消息、topic 16 分区时,sticky 的 p99 延迟约为默认(逐条轮询)策略的一半。
- 分区越多,优势越大:固定 3 producer、每秒 10000 条,把分区数从 16 → 64 → 128,默认策略的延迟增长快得多;即便只有 16 分区,默认策略的平均 p99 也是 sticky 的约 1.5 倍。
linger.ms > 0时收益更夸张:低吞吐 +linger.ms=1000场景下(1 producer、1000 msg/s),默认策略 p99 是 sticky 的约 5 倍——因为 sticky 能更快把批填满、不用死等 linger 到点。- 常常还降 CPU:3 producer、10000 msg/s、16 分区时,sticky 明显降低了 CPU 占用(更少批次 = 更少开销)。
5.3 有 key / 混合 key:三个”零副作用”验证
Confluent 特意测了三种带 key 的场景,来证明 sticky 那点额外逻辑不会拖累有 key 路径:
| 场景 | 结果 | 原因 |
|---|---|---|
| 随机 key + 无 key 混合 | 批处理略有改善,但整体收益不显著 | 有 key 消息走哈希、忽略 sticky |
| 纯随机 key(高吞吐) | 与默认基本持平 | 无 sticky 行为、无额外批处理,多出的逻辑不影响延迟 |
| 顺序 key + 大量分区 | 无明显差异(这是最”不利”的场景,几乎每条都新建批次) | 额外逻辑集中在新建批次附近,实测也没有引入可见开销 |
一句话总结基准:sticky 在”无 key + 分区多 + linger 有等待”时收益最大,在有 key 场景零副作用。 这就是它能在 2.4 直接成为默认行为的底气。
6. 自定义分区器:把热点 key 单独隔离
大多数团队用不上自定义分区器,但它是应对热点 key 倾斜的正解之一。经典场景:一张交易流水表,用户 CEO 一个人贡献了 40% 以上的交易——用默认哈希,他会和别的用户挤在同一个分区,把那个分区撑爆、拖慢对应消费者。
思路是给热点 key 一个专属分区,其余 key 照常哈希到剩下的分区:
public class HotKeyPartitioner implements Partitioner {
@Override
public int partition(String topic, Object key, byte[] keyBytes,
Object value, byte[] valueBytes, Cluster cluster) {
int numPartitions = cluster.partitionsForTopic(topic).size();
// 约定:分区 0 专供热点 key,其余 key 哈希到 [1, numPartitions)
if ("CEO".equals(key)) {
return 0;
}
int hash = Utils.toPositive(Utils.murmur2(keyBytes));
return 1 + hash % (numPartitions - 1);
}
@Override public void close() {}
@Override public void configure(Map<String, ?> configs) {}
}
partitioner.class=com.example.HotKeyPartitioner
代价与取舍:隔离热点保住了其余分区的均衡与消费吞吐,但专属分区自身仍是单点热点——它的上限就是单分区上限。真正超大的热点还得叠加下游扩容、二级 sharding 或采样等手段。自定义分区器”不复杂但也不轻量”,只在标准策略确实覆盖不了时才上。
7. 顺序性:分区策略绕不开的另一面
分区策略的另一个后果是顺序性。Kafka 的铁律(见核心原理篇):只保证单个 partition 内有序,不保证跨分区全局有序。
于是顺序性完全由”消息去了哪个分区”决定,也就是又回到了分区策略:
7.1 要顺序,就用 key 把相关消息钉到同一分区
典型场景:先”改用户会员等级”、再”按等级算订单价”,两条消息顺序错了结果就错。做法:用业务标识(订单 ID / 用户 ID)作为 key,相同 key 经哈希必落同一分区,天然有序。
// 默认分区器:相同 key 的消息进同一分区 → 分区内有序
producer.send(new ProducerRecord<>("order-topic", orderId, orderMsg));
消费侧只要保证同一 group 内一个分区只被一个消费者消费(这是 Kafka 天然保证的),按接收顺序处理即可。注意 §2.1 的前提:这条保证以”分区数不变”为条件。
7.2 “全局顺序”的代价
想要全局严格有序,唯一彻底的办法是整个 topic 只用 1 个分区——但这直接放弃了 Kafka 的并行与扩展,吞吐塌方。所以实践里几乎总是选择**“局部有序”**:按 key 保证每个业务实体(每个订单/每个用户)内部有序,不同实体之间不要求顺序。
7.3 多分区下”既要顺序又要吞吐”
常见工程手法:消费者先按分区顺序把消息塞进内存队列(如 BlockingQueue),再用线程池按入队顺序处理——在保住”相同业务标识局部有序”的同时,用异步/多线程把消费端吞吐提上去(代价是要处理好队列溢出、线程同步、失败重试)。
7.4 严格有序还要防”重试乱序”
有个隐藏坑:Producer 重试可能让同一分区内的消息乱序(前一批失败重发、后一批已成功)。要严格有序:
enable.idempotence = true # 幂等生产者(推荐,自动处理下面两项)
max.in.flight.requests.per.connection <= 5 # 幂等下可保序
retries > 0
开启幂等后,Kafka 会给消息编序号,即使重试也能保证分区内不重复、不乱序。这也是为什么”顺序性”这件事最终还是绕回到 Producer 的配置。
8. 怎么选:一张决策表
把生产侧分区策略的旋钮合到一起(消费侧 assignor 的选择见消费者篇 §4.1):
| 你的诉求 | 选择 |
|---|---|
| 无 key、只要吞吐/低延迟 | 默认(2.4+ = Sticky),什么都不用改 |
| 需要”相同业务实体有序” | 用 key(订单/用户 ID),走默认哈希;并把分区数当契约、别随意扩 |
| 需要严格不乱序(防重试乱序) | 用 key + enable.idempotence=true + 限 max.in.flight |
| 单个热点 key 压垮某分区(且不需按该 key 保序) | Round-robin 或 自定义分区器隔离热点 |
| 想手工把某类消息钉到指定分区 | 自定义 Partitioner |
| 全局严格有序 | 单分区(谨慎,吞吐代价大,尽量改成”按 key 局部有序”) |
对绝大多数团队:无 key 场景什么都别调,享受 2.4 默认的 sticky;有顺序诉求就上 key + 幂等。 这两条覆盖了写入侧 90% 的需求;消费侧的负载均衡与平滑重平衡,统一上
CooperativeStickyAssignor(见消费者篇 §4.3)。
9. 参数速查
| 参数 | 位置 | 默认 | 作用 |
|---|---|---|---|
partitioner.class | Producer | 默认分区器 | 选择/自定义分区策略 |
batch.size | Producer | 16384 B | 单批写满即发的字节阈值 |
linger.ms | Producer | 0 | 批未满时的最长等待;调大能进一步增大批、放大 sticky 收益 |
enable.idempotence | Producer | true(≥3.0) | 幂等:分区内不重复、不乱序 |
max.in.flight.requests.per.connection | Producer | 5 | 单连接未确认请求上限;幂等下 ≤5 可保序 |
acks | Producer | all(≥3.0) | 确认级别,影响”不丢”与延迟 |
消费侧的
partition.assignment.strategy/group.instance.id参数属于消费链路,见消费者篇 §4。
10. 落到 AdTech:为什么这事对广告链路很实在
广告事件流(曝光/点击/转化)常常是无 key 的高吞吐洪峰——正是 sticky partitioner 的最佳受益场景:更大的批、更少的请求,直接压低事件入 Kafka 的延迟与 broker CPU。
而需要顺序的地方也很典型:同一次曝光的”曝光 → 点击 → 转化”要能按序归因,就用 请求 ID / 用户 ID 作为 key 把同一实体的事件钉到同一分区,保证归因链路的局部有序——既拿到顺序,又不牺牲整体并行。同一个分区策略旋钮,两种诉求各取所需。
至于个别超大广告主造成的 key 倾斜,则用轮询或自定义分区器把它从”拖垮一个分区”变成”摊平到全体”。而消费侧多条链路(计费、特征、风控)各自成组、怎么用 CooperativeStickyAssignor 扛住发布期的 rebalance,属于消费链路的取舍,见消费者篇 §10。
11. 常见误解 ↔ 正解
| 常见误解 | 正解 |
|---|---|
linger.ms=0 就是每条单独发 | 不是。同时到达同一分区的消息仍会被自然攒批 |
| 轮询最均匀所以最好 | 逐条轮询造成大量小批次 → 更高延迟;sticky 长期同样均匀却批更大 |
| 轮询一无是处 | 它是”有 key 但严重倾斜”时的负载均衡解药(代价是放弃按 key 保序) |
| Sticky 会让数据分布不均 | 短期集中、长期均匀(不断换分区),分布基本不受影响 |
| Sticky 对所有消息都提速 | 只优化无 key消息;有 key 走哈希、无收益也无副作用 |
| 用了 key 就永远同分区、永远有序 | 只在分区数不变时成立;扩分区会改变取模结果、打破历史顺序 |
| 加分区就能提高顺序性 | 恰恰相反,分区越多越难全局有序;顺序靠 key 钉到同一分区 |
| 用了 key 就一定不乱序 | 还需 enable.idempotence + 限制 in-flight,否则重试可能乱序 |
| Sticky 要手动开启 | Kafka 2.4+ 默认分区器已内置,无需配置 |
| ”分区策略”只是生产者的事 | 消费侧的 assignor(分区分给哪个消费者)同样关键,决定负载均衡与 rebalance 代价(见消费者篇 §4.1) |
12. 速查表
- 发送 = 攒批:
RecordAccumulator按分区攒批,batch.size满或linger.ms到才发。批大批少 → 延迟低。 - 能否攒批取决于分区:只有同分区消息能进同一 batch。
- 四策略:默认(有 key 哈希/无 key 走 sticky)、轮询、Sticky、自定义。
- Sticky(2.4+):无 key 消息粘住一分区灌满再换 → 批更大、请求更少、p99 延迟约砍半、常降 CPU;对有 key 零副作用。
- key 路由前提:哈希只在分区数不变时稳定;扩分区会打破”同 key 同分区”。
- 热点 key:不需保序 → 轮询;需要隔离 → 自定义分区器给专属分区。
- 顺序:只保证分区内有序;要顺序 → 用 key;要严格不乱序 → key +
enable.idempotence+ 限 in-flight;全局有序 → 单分区(慎用)。 - 默认建议:无 key 不用调(享受 sticky);有顺序诉求上 key + 幂等。
- 另一半在消费侧:分区”分给哪个消费者”由 assignor 决定(推荐 CooperativeSticky),拆解见消费者篇 §4。
分区策略是 Kafka Producer 上”最小改动、最大杠杆”的旋钮之一:Sticky Partitioner 只是换了个”先灌满再换”的分配顺序,就把无 key 高吞吐场景的延迟砍了一半。理解它,再把消费侧的 assignor(消费者篇 §4)一并想清楚,也就理解了 Kafka”吞吐 / 延迟 / 顺序 / 均衡四件事都由分区决定”的底层逻辑。
延伸阅读
本系列内部串读(Kafka 深挖系列):
- Kafka 核心原理精讲:从一条日志到分布式流平台:Partition / Offset / Consumer Group / Rebalance 与”只保证分区内有序”的由来。
- Kafka Consumer 深挖:poll 循环、offset 提交与 rebalance:写入侧的对称另一半——消费侧的分区分配(assignor)、提交时机、rebalance 与并行模型。
- Kafka 可靠性与 Exactly-Once 深挖:
acks/enable.idempotence/max.in.flight背后的复制与事务机制。 - Kafka 性能内核:磁盘系统凭什么跑出内存级吞吐:批量与压缩是怎样和顺序写、零拷贝一起放大吞吐的。
- 为什么 AdTech 偏爱 Kafka:广告事件中枢的五大典型场景与架构设计:无 key 洪流、热点隔离在广告采集/计费/风控链路里的落地。
一手资料与优质教程:
- Confluent. Apache Kafka Producer Improvements: Sticky Partitioner:Sticky Partitioner 的官方设计、基准方法(Trogdor/Castle ProduceBench)与数据来源(Justine Olshan, KIP-480)。
- Redpanda. Kafka Partition Strategy:四种生产者分区策略 + 消费者分配策略(Range/RoundRobin/Sticky/Custom)与 rebalance 的系统梳理。
- Apache Kafka. Producer Configs:
batch.size/linger.ms/partitioner.class/enable.idempotence的权威参数说明。