Skip to content
Charles Shao
Go back

Kafka Producer 深挖:分区策略与 Sticky Partitioner 是怎么把延迟砍半的

views

前一篇立了地基:Kafka 是一份分布式可重放日志。这一篇聚焦在 Producer 端一个看起来微不足道、实则牵一发动全身的决策:

一条消息,到底该写进哪个 partition?

这个决策(分区策略)同时决定了吞吐、延迟、顺序性。而 Kafka 2.4 引入的 Sticky Partitioner 就是一次针对这个决策的精妙优化——在几乎不改变数据分布的前提下,把无 key 消息的 p99 延迟砍掉约一半。下面把来龙去脉讲透。(分区写进去后”分给哪个消费者”的另一半——消费侧 assignor——放在消费者篇 §4 讲。)

TL;DR

Table of contents

Open Table of contents

1. 先看生产者发送流程:一切从”攒批”说起

要理解分区策略为什么重要,得先知道 Producer 不是收到一条就发一条。先把要优化的指标定义清楚:

生产者延迟 ≈ 一条消息从被客户端 send() 到**被 Kafka 确认(acknowledged)**所经历的时间。降低这个时间,整条链路都受益——而它的大头,往往花在”攒批与请求排队”上。

一条消息的旅程:

  1. 封装成 ProducerRecord序列化
  2. 交给 Partitioner 算出目标分区;
  3. 写进 RecordAccumulator,按分区攒成 RecordBatch
  4. 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 消息的方式是逐条轮询——来一条,发一个分区;再来一条,发下一个分区……看起来最公平,但它和”攒批”是天生的冤家:

Round-robin 逐条轮询 vs Sticky 黏性分区的批处理对比图。上半部分『Round-robin(2.4 前)』:Producer 依次产生 6 条无 key 消息 m1m6,被逐条轮流分配到 3 个分区——m1→P0、m2→P1、m3→P2、m4→P0、m5→P1、m6→P2。结果每个分区只攒到 2 条消息,形成 3 个很小的批次(每批 2 条),于是要发 3 个请求,每条消息分摊的批处理开销高、排队多、延迟高。下半部分『Sticky(2.4 后)』:同样 6 条消息 m1m6,Producer 先『粘住』分区 P0 连续发 m1~m4 直到该批次写满/完成,再随机换到 P1 发 m5、m6。结果 P0 攒成一个较大的批次(4 条),P1 一个批次(2 条),P2 本轮为空——批次更少更大,只需更少的请求,单条成本更低、延迟更低。图侧标注:长期来看多轮之后 sticky 依旧让数据在各分区间大致均匀,但每一轮都换来了更优的批处理。

问题就在这张图的上半部分:逐条轮询把消息摊得到处都是,每个分区同一时刻只攒到寥寥几条,于是形成一大堆小批次。 而:

小批次 = 更多请求 + 更多排队 = 更高延迟。 每处理一个 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 factor3
Topic 数4
linger.ms0(除单独说明的 linger 实验外)
acksall
key 生成全部为 null(考察无 key 路径)
useConfiguredPartitioner / skipFlushtrue / true

两个 true 是关键:useConfiguredPartitioner 确保用 DefaultPartitioner 真正决定分区,skipFlush 确保批次是靠写满或 linger.ms 触发而不是被 flush 逼出——这样对比的才是”分区策略本身”的效果。

5.2 核心结论

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.classProducer默认分区器选择/自定义分区策略
batch.sizeProducer16384 B单批写满即发的字节阈值
linger.msProducer0批未满时的最长等待;调大能进一步增大批、放大 sticky 收益
enable.idempotenceProducertrue(≥3.0)幂等:分区内不重复、不乱序
max.in.flight.requests.per.connectionProducer5单连接未确认请求上限;幂等下 ≤5 可保序
acksProducerall(≥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. 速查表

分区策略是 Kafka Producer 上”最小改动、最大杠杆”的旋钮之一:Sticky Partitioner 只是换了个”先灌满再换”的分配顺序,就把无 key 高吞吐场景的延迟砍了一半。理解它,再把消费侧的 assignor(消费者篇 §4)一并想清楚,也就理解了 Kafka”吞吐 / 延迟 / 顺序 / 均衡四件事都由分区决定”的底层逻辑。


延伸阅读

本系列内部串读(Kafka 深挖系列):

一手资料与优质教程:


views
Share this post on:

Previous Post
Kafka Consumer 深挖:poll 循环、offset 提交与 rebalance,怎么做到不丢不重不卡顿
Next Post
Kafka 核心原理精讲:从一条日志到分布式流平台,把关键机制与取舍讲透