Kafka Producer 完整执行链:一条 Record 从 send() 到 Partition Leader 的旅程
从 Metadata、Serializer、Partition 选择、RecordAccumulator 到 Sender 与 NetworkClient
阅读目标
第一篇建立了 Kafka 的 Log → Partition → Offset → Consumer Position 世界观;第二篇继续把 Partition 放大,解释了 RecordBatch、LogSegment、稀疏索引、Page Cache 与 Broker 侧读写路径。
但前两篇都从一个已经形成的 ProduceRequest 开始观察。
现在把视角移回 Java 应用:
Future<RecordMetadata> future = producer.send(record);这一行代码执行以后,消息是否已经离开 JVM?
答案通常是否定的。
send() 的主要职责不是在业务线程中完成一次网络往返,而是把一条业务 Record 转换成可进入 Kafka Client Runtime 的待发送状态。真正的批次调度、连接管理、请求发送、重试推进和响应完成,由 Producer 内部的后台 I/O 线程持续执行。
完成本篇后,你应该能够:
- 画出 Kafka Producer 的应用线程与 Sender 线程模型;
- 从
KafkaProducer.send()追踪到RecordAccumulator.append(); - 解释
bootstrap.servers为什么不是完整 Broker 路由表; - 区分 Metadata、Topic Partition 信息与 Partition Leader;
- 解释显式 Partition、有 Key 和无 Key 时的路由逻辑;
- 解释为什么 RecordAccumulator 按
TopicPartition组织 Deque; - 区分 Record、ProducerBatch、ProduceRequest 三种聚合层次;
- 解释 BufferPool 如何把 Broker 或网络背压传导到业务线程;
- 画出 Sender 的
ready → drain → request → poll → complete循环; - 区分
max.block.ms、linger.ms、request.timeout.ms与delivery.timeout.ms; - 根据 Producer Metrics 判断瓶颈位于 Metadata、Buffer、Batch、Network 还是 Broker;
- 不再把 Kafka Producer 简化成“序列化后通过 Socket 发出去”。
本文以 Apache Kafka 4.3 官方文档和 4.x Java Client 源码结构为主线。Kafka 4.3 中 linger.ms 默认值为 5 ms,batch.size 默认 16 KiB,buffer.memory 默认约 32 MiB,幂等生产在无冲突配置时默认启用;源码阅读时仍应固定到实际使用的 Kafka Tag,而不是依赖持续变化的 trunk 行号。13
0. send() 返回了,消息真的发出去了吗
看一段非常普通的代码:
ProducerRecord<String, String> record = new ProducerRecord<>("order-events", "order-1001", json);
Future<RecordMetadata> future = producer.send(record);
System.out.println("send returned");很多人的直觉是:
send() ↓建立网络连接 ↓把消息写入 Broker ↓Broker 返回 ACK ↓方法返回如果真是这样,每条 Record 都要在业务线程里完成一次网络往返,Kafka Producer 很难形成高吞吐。
实际主模型更接近:
Application Thread ↓KafkaProducer.send() ↓Serialize / Partition ↓RecordAccumulator.append() ↓返回 Future<RecordMetadata>
================ Thread Boundary ================
Sender Thread ↓选择可发送 Batch ↓构建 ProduceRequest ↓NetworkClient ↓Partition Leader ↓收到响应 ↓完成 Future / Callback
Kafka 官方 API 将 send() 定义为异步操作:Record 进入待发送 Buffer 后,调用可以返回,从而允许 Producer 将多条 Record 合并成更高效的批次。2
但要立刻补一个重要边界:
异步发送不等于
send()永远不会阻塞。
调用线程至少可能在以下位置停留:
- Topic Metadata 尚不可用,需要等待 Metadata 更新;
- BufferPool 内存不足,需要等待 Sender 释放 Buffer;
- 用户自定义 Interceptor、Serializer 或 Partitioner 本身执行很慢;
- Record 过大或配置非法,可能在调用线程直接失败;
- Producer 已经关闭或进入不可恢复状态,可能同步抛出异常。
max.block.ms 限制的是 send() 等待 Metadata 和 Buffer 分配的累计时间,但用户 Serializer 和 Partitioner 的执行时间并不计入这一超时。1
因此,判断一条消息的状态时需要区分:
send() 已返回 ≠网络请求已发送 ≠Broker 已接收 ≠Record 已写入 Leader Log ≠满足目标 ACK 语义这些状态之间通过 Producer 内部的异步 Runtime 串联起来。
1. Kafka Producer 的整体对象与线程模型
1.1 KafkaProducer 不是一个轻量方法包装器
创建一个 KafkaProducer 时,客户端会构造并持有一组长期运行的组件:
KafkaProducer├── ProducerConfig├── Serializer<Key>├── Serializer<Value>├── ProducerInterceptors├── ProducerMetadata├── RecordAccumulator├── BufferPool├── Sender├── NetworkClient / KafkaClient├── Metrics└── I/O Thread因此,Producer 的正确使用方式通常是:
// 作为长期复用、线程安全的客户端实例KafkaProducer<String, OrderEvent> producer = new KafkaProducer<>(props);而不是每发送一条 Record 就:
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) { producer.send(record);}频繁创建 Producer 会重复执行:
- 配置解析;
- Serializer 与 Interceptor 初始化;
- Metadata Bootstrap;
- 网络连接建立;
- Sender 线程创建;
- BufferPool 分配;
- 指标注册与销毁。
KafkaProducer 是线程安全的,多个业务线程共享一个实例通常比每个线程创建一个实例更合理。2
1.2 两类线程承担不同职责
Producer 的核心不是“多线程越多越快”,而是把两种职责分开。
Application Threads
执行:
创建 ProducerRecord调用 Interceptor等待 MetadataSerializerPartition 选择Record Size 校验RecordAccumulator Append返回 FutureSender I/O Thread
执行:
检查 Batch 是否 Ready处理过期 Batch刷新 Metadata按 Broker Node Drain Batch构建 ProduceRequest管理连接发送请求poll 网络事件处理响应或错误完成 Callback / Future
线程边界带来三个关键能力。
第一,批处理
多个业务线程调用 send() 时,Record 可以先聚集到按 Partition 组织的 Batch 中,再由一个 Sender 统一发送。
第二,调用与网络解耦
业务线程不需要为每条 Record 执行一次完整 Broker 往返。
第三,显式背压边界
如果 Sender 长期发送不过去,BufferPool 会被消耗,最终使新的业务线程在 send() 内等待甚至失败。
这说明 Producer 并没有“消灭背压”,只是把背压从同步 RPC 响应时间转换为:
Accumulator Queue DepthBufferPool Available MemoryBatch AgeRequest LatencyRetry / Timeout1.3 Producer 的资源生命周期
Producer 的生命周期通常包括:
Construct ↓Bootstrap Metadata ↓Repeated send() ↓Background Sender Loop ↓flush() 可选等待 ↓close()send()
把 Record 登记进异步发送链。
flush()
等待当前已发送或已登记的 Record 完成,但会牺牲异步流水线能力。不要在每条 Record 后调用。
close()
停止新发送,等待或放弃剩余请求,关闭 Sender、网络连接、Serializer、Interceptor 和 Metrics。
如果应用直接终止而没有合理关闭,Accumulator 中尚未发送或尚未完成的 Record 可能丢失。
2. KafkaProducer.send() 的完整主路径
2.1 先看压缩后的执行链
现代 Kafka Client 的具体方法签名会演进,但理解 Producer 的稳定主链可以从下面开始:
KafkaProducer.send(record) ↓ProducerInterceptors.onSend ↓doSend ↓等待 / 获取 Topic Metadata ↓Serialize Key / Value ↓选择 Partition ↓估算 Record Size 并校验 ↓RecordAccumulator.append ↓必要时唤醒 Sender ↓返回 Future<RecordMetadata>
这里最重要的认知是:
send()的主路径结束于“Record 已进入 Client Buffer”,而不是“Broker 已完成持久化”。
2.2 ProducerRecord 只是业务输入对象
典型 Record:
ProducerRecord<String, OrderEvent> record = new ProducerRecord<>( "order-events", null, // partition System.currentTimeMillis(), "order-1001", // key event, // value headers);它可能包含:
TopicOptional PartitionOptional TimestampOptional KeyValueHeaders此时 Key 和 Value 仍然可能是 Java 对象,尚未变成 Kafka 协议中的 byte[]。
2.3 Interceptor 在序列化前获得 Record
ProducerInterceptor.onSend() 可以:
- 增加 Header;
- 注入 Trace ID;
- 记录审计字段;
- 修改 Topic、Key 或 Value;
- 收集自定义指标。
但 Interceptor 运行在调用线程中。
因此下面的做法很危险:
public ProducerRecord<String, String> onSend(ProducerRecord<String, String> record) { remoteAuditService.report(record); // 同步网络调用 return record;}它会把 Producer 的异步发送入口重新变成同步依赖。
Interceptor 还必须注意:
- 不要修改不应变化的业务语义;
- 不要执行高延迟 I/O;
- 不要抛出不可控异常;
- 不要把大对象复制到 Header;
- 不要在
onAcknowledgement中执行耗时逻辑。
官方 ProducerInterceptor 接口允许在发布前拦截并修改 Record,但插件代码仍属于应用自身的延迟与稳定性责任。4
2.4 Metadata 为什么必须先于 Partition 选择
Producer 要把 Record 发送到 Topic 的某个 Partition Leader,至少需要知道:
Topic 是否存在?Topic 有多少 Partition?每个 Partition 当前 Leader 是谁?Leader 对应哪个 Broker Node?如果连 Partition 数量都不知道,就无法正确执行:
hash(key) % partitionCount因此 send() 可能需要等待目标 Topic 的 Metadata 可用。
如果等待超过 max.block.ms,调用可能失败;这也是为什么 Kafka 集群故障时,业务线程可能在 send() 入口积压,而不是一直“无感异步”。
2.5 Serializer 把对象转换为协议字节
Kafka Broker 不理解 Java 的 OrderEvent 对象。
Producer 必须调用:
keySerializer.serialize(...)valueSerializer.serialize(...)得到:
serializedKey: byte[]serializedValue: byte[]Serializer 的职责应当尽量单纯:
Java Object ↓Stable Binary Representation不建议在 Serializer 中:
- 查询数据库补字段;
- 调用远程 Schema 服务且无本地缓存;
- 修改业务状态;
- 写外部审计表;
- 执行不可控重试;
- 依赖线程本地脆弱上下文。
因为这些逻辑发生在业务调用线程,并且其耗时不受 max.block.ms 约束。1
2.6 Record Size 校验发生在入队前
序列化后,Producer 已经知道 Key、Value 和 Header 的字节规模,可以检查:
Record 是否超过 max.request.size?Record 是否超过 Producer 可接受的 Batch / Request 边界?Header 是否异常膨胀?Kafka 4.3 中 max.request.size 默认 1 MiB,它既限制单个请求规模,也间接限制未压缩 RecordBatch 的上界;Broker 端仍有独立消息大小限制,二者需要协同配置。1
专家级排查不能只改 Producer:
Producer max.request.sizeBroker message.max.bytesTopic max.message.bytesReplica fetch sizeConsumer fetch size任何一个链路边界不兼容,都可能导致“大消息在某一层失败”。
2.7 RecordAccumulator.append 是调用线程主链终点
完成 Metadata、序列化、Partition 选择和大小校验后,Producer 尝试把 Record 追加到目标 TopicPartition 的当前 Batch。
结果通常包含:
- 对应 Record 的 Future;
- 是否创建了新 Batch;
- 当前 Batch 是否已经 Full;
- 是否需要唤醒 Sender;
- 是否需要重新执行 Partition 选择或处理关闭状态。
如果 Batch 已满或刚创建了新 Batch,Producer 会唤醒 Sender,使它无需一直等待下一次定时 poll。
3. Metadata:Producer 如何找到 Partition Leader
3.1 bootstrap.servers 只负责 Bootstrap
配置:
bootstrap.servers=kafka-1:9092,kafka-2:9092,kafka-3:9092并不意味着 Producer 永远只和这三个地址通信,也不要求列出集群全部 Broker。
它的作用是:
提供一组初始可连接地址,让 Client 进入 Kafka Cluster 并获取完整 Metadata。
官方建议配置多个 Bootstrap 地址提高初始连接韧性,但 Client 会在 Bootstrap 后发现并管理完整 Broker 集合。1

3.2 Producer 为什么可以直连 Leader
Kafka 不要求所有 Produce 流量经过一个中央 Router。
Producer 获取 Metadata 后,可以建立映射:
Topic: order-eventsP0 Leader → Broker 1P1 Leader → Broker 3P2 Leader → Broker 2然后:
orders-P1 Batch ↓直接发送 Broker 3Kafka 官方设计明确说明,Producer 直接向目标 Partition Leader 发送数据;集群节点可以回答 Metadata 请求,使 Client 知道 Broker 与 Partition Leader 的位置。5
这种设计减少了中央路由层,但也把一部分复杂性放到 Client:
- Metadata Cache;
- Leader 变化感知;
- 连接池;
- 按 Broker 聚合请求;
- 错误后刷新路由。
3.3 Metadata 缓存保存什么
简化理解:
Cluster Metadata├── Brokers├── Topics│ ├── Partition Count│ └── Partition Metadata│ ├── Leader│ ├── Replicas│ └── ISR 摘要└── Cluster / Version StateProducer 最关注:
TopicPartition → Leader Node副本与 ISR 的完整可靠性语义放在第五篇。
3.4 Metadata 何时刷新
常见触发条件包括:
- 第一次发送某个 Topic;
- Metadata 缺失;
- Leader 已变化;
- Broker 返回需要更新 Metadata 的错误;
- Topic Partition 数发生变化;
metadata.max.age.ms到期;- Topic 长期空闲后缓存被忘记;
- 已知 Broker 全部不可用,需要 Rebootstrap。
Kafka 4.3 支持在 Client 无法从已知 Broker 获取 Metadata 时,按配置重新使用 bootstrap.servers 进行 Rebootstrap。1
3.5 Metadata 等待为什么会卡住 send()
假设应用第一次发送:
order-eventsProducer Metadata 尚不知道这个 Topic。
Application Thread ↓waitOnMetadata(order-events) ↓Sender 被唤醒 ↓MetadataRequest ↓Broker Response ↓Application Thread 获得 Partition 信息所以 send() 虽然是异步发送 API,但第一次发送、新 Topic、Leader 抖动或 Cluster 不可达时,调用线程仍可能等待。
3.6 Metadata 问题的典型症状
症状一:首次发送延迟明显
可能是正常 Bootstrap、DNS、TLS 握手或 Topic Metadata 获取。
症状二:大量 TimeoutException: Topic ... not present in metadata
可能原因:
- Topic 不存在且禁止自动创建;
- ACL 不允许 Describe;
- Advertised Listener 不可达;
- DNS 解析错误;
- Broker 返回的地址对应用网络不可见;
- Metadata 请求一直失败。
症状三:Leader 切换期间发送抖动
Producer 需要接收错误、标记 Metadata 过期、刷新 Leader,再重试可恢复 Batch。
不要一看到发送超时就直接增加 delivery.timeout.ms。如果根因是错误的 advertised.listeners,延长超时只是让失败更慢。
4. Serializer 与 Partition 选择
4.1 Partition 决定顺序域和并行域
第一篇已经建立:Kafka 只保证 Partition 内顺序。
因此 Producer 的 Partition 选择同时决定:
哪些 Record 进入同一有序 Log哪些 Record 可以并行写入哪些 Consumer 后续可以并行处理数据是否形成 Hot PartitionPartitioner 不是一个“负载均衡小工具”,而是业务语义到分布式 Log 的映射函数。
4.2 默认选择顺序
现代 Kafka 默认逻辑可以概括为:
1. Record 显式指定 Partition → 使用显式 Partition
2. 未指定 Partition,但存在 Key → 根据 Key Hash 选择 Partition
3. 既没有 Partition,也没有 Key → 使用 Sticky / Adaptive 选择策略形成更有效 Batch
Kafka 4.3 官方配置说明指出:默认策略在有 Key 时基于 Key Hash 选择 Partition;无 Key 时会使用 Sticky Partition,并在至少形成约 batch.size 数据后切换。自适应 Partitioning 默认开启,可结合 Broker 处理表现调整无 Key 数据的分布。1
4.3 显式 Partition
new ProducerRecord<>("order-events", 2, key, value);优点:
- 精确控制路由;
- 可实现特殊兼容策略;
- 便于某些固定 Shard 设计。
风险:
- Partition 扩容后代码可能过时;
- 指定不存在 Partition 会失败;
- 容易产生热点;
- 业务代码与 Topic 物理布局耦合;
- 多应用必须保持同一规则。
除非有清晰架构理由,不要把 Partition 编号硬编码进普通业务逻辑。
4.4 有 Key:局部顺序与数据局部性
例如:
Key = orderIdOrderCreated(order-1001)OrderPaid(order-1001)OrderShipped(order-1001)只要 Key 的序列化结果、Partition 数和分区算法保持一致,这些 Record 会进入同一 Partition,从而获得局部顺序。
需要注意:
Hash 输入是序列化后的 Key Bytes,而不是 Java 对象的
hashCode()。
自定义 Key Serializer 变化可能改变路由结果。
4.5 无 Key:为什么不是简单逐条 Round Robin
如果每条无 Key Record 都切换 Partition:
R1 → P0R2 → P1R3 → P2R4 → P0每个 Partition 的 Batch 可能都很小。
Sticky 思路是:
一段时间集中写 P1 ↓形成较满 Batch ↓切换到其他 Partition目标不是保证每一瞬间绝对均匀,而是改善:
- Batch 填充率;
- 压缩率;
- Request 数量;
- Broker 与 Client I/O 摊销。
4.6 自定义 Partitioner 的边界
自定义 Partitioner 适用于:
- 兼容已有 Sharding 规则;
- 特定租户隔离;
- 地域或业务域局部性;
- 特定 Hot Key 拆分协议。
但必须回答:
Partition 扩容怎么办?规则版本如何发布?多个 Producer 是否保持一致?Key 是否稳定?热点如何检测?旧消息和新消息的顺序假设是否仍成立?而且 Partitioner 运行在业务调用线程,不能执行远程 I/O。
5. RecordAccumulator:Producer 的核心缓冲结构
5.1 为什么不能直接把 Record 交给 NetworkClient
假设每次调用:
producer.send(record);都立即创建请求:
Record ↓Request Header + Protocol Encoding ↓Socket Write大量小 Record 会导致:
- 请求数过多;
- 系统调用频繁;
- TCP 小包增多;
- Broker 请求处理开销增大;
- 压缩效果差;
- 磁盘追加 Batch 过小。
因此 Producer 在业务线程与 NetworkClient 之间引入:
RecordAccumulator它把单条 Record 转换为按 Partition 聚集的 Batch。
5.2 为什么按 TopicPartition 组织 Deque
RecordAccumulator 的概念结构可以表示为:
Map<TopicPartition, Deque<ProducerBatch>>
例如:
orders-0 → [Batch A][Batch D]orders-1 → [Batch B]orders-2 → [Batch C][Batch E]为什么不能把不同 Partition 混进同一个 ProducerBatch?
因为一个 Batch 需要共享:
- TopicPartition;
- Partition Leader 路由;
- Offset 分配域;
- 顺序域;
- 压缩与 RecordBatch Header;
- 幂等 Sequence 范围;
- Broker 端 Append 目标。
因此:
ProducerBatch 的聚合边界是 TopicPartition。
但一次 ProduceRequest 可以向同一个 Broker 携带多个 Partition 的 Batch。
5.3 Append 的两种主要路径
路径一:追加到现有 Batch
找到 TopicPartition Deque 尾部 Batch ↓剩余空间足够 ↓把 Record 编码进 Batch ↓返回 Future路径二:创建新 Batch
没有 Batch 或当前 Batch 空间不足 ↓向 BufferPool 申请内存 ↓创建 MemoryRecordsBuilder / ProducerBatch ↓加入 Deque 尾部 ↓追加 Record ↓必要时唤醒 Sender真实源码为了并发、关闭、分区重新选择和事务状态包含更多分支,但核心矛盾不变:
能否复用当前 Batch?否则能否申请新 Buffer?5.4 batch.size 是预期 Batch Buffer 上界,不是发送请求大小
Kafka 4.3 默认:
batch.size = 16384 bytes它表示 Producer 为每个待创建 Batch 使用的默认 Buffer 规模。1
不要把它理解为:
每个 ProduceRequest 最大 16 KiB因为一次 Request 可以包含:
Broker Node 1├── orders-0 Batch├── payments-2 Batch└── audit-7 BatchRequest 总上限更多受 max.request.size 约束。
5.5 大于 batch.size 的单条 Record 怎么办
官方配置说明指出,Producer 不会尝试把大于 batch.size 的 Record 与其他 Record 按普通方式合批,但这不意味着它一定被 batch.size 拒绝。Client 可以为单个大 Record 分配足够容纳它的 Buffer,只要没有超过更高层的消息和请求大小限制。1
所以:
batch.size≠单条消息绝对最大值真正需要联合检查:
estimated serialized sizemax.request.sizeBroker / Topic message limitsBuffer availability5.6 linger.ms 控制什么
Batch 未满时,如果立即发送,可能形成大量小 Batch。
linger.ms 给 Sender 一个批量等待上界:
Batch 已满 → 立即 Ready
Batch 未满但等待时间达到 linger.ms → Ready
Broker Backpressure / Connection Not Ready → 实际等待可能更久Kafka 4.0 起,linger.ms 默认从 0 调整为 5 ms,官方解释是更大 Batch 带来的效率提升通常可以抵消甚至改善实际延迟。1
注意:
linger.ms=5不代表每条 Record 必定等待 5 ms。
高负载下 Batch 很快填满,会提前发送;Sender 的 Request 节奏也会自然聚集期间到达的 Record。
5.7 压缩在哪里发生
压缩以 Batch 为边界。
Record 1Record 2Record 3 ↓RecordBatch Compression更合理的 Batch 往往提高压缩率,因为业务消息之间通常共享大量字段和值模式。Kafka 官方配置明确说明,压缩针对完整 Batch,Batching 效率会影响压缩率。1
压缩带来的 Trade-off:
Network Bytes ↓Disk Bytes ↓Broker I/O ↓
Producer CPU ↑Consumer CPU ↑Latency 可能变化不能脱离 CPU、网络与消息类型直接宣称某个 Codec 永远最优。
5.8 Batch 是如何与 Future 关联的
一次 ProducerBatch 内可以包含多个 Record,每个 Record 调用 send() 都需要自己的完成结果:
Batch├── Record A → Future A / Callback A├── Record B → Future B / Callback B└── Record C → Future C / Callback CBroker 对 Batch 返回 Base Offset、Timestamp 和 Error 后,Producer 再根据 Record 在 Batch 中的位置完成各自的 RecordMetadata。
因此 Future 的完成由 Sender 处理响应触发,不由 send() 调用线程完成。
6. BufferPool:内存限制与背压传播
6.1 buffer.memory 限制的是什么
Kafka 4.3 默认:
buffer.memory = 33554432 bytes约 32 MiB。
它近似限制 Producer 可用于缓存待发送 Record 的总内存,但官方特别说明它不是 Producer 总内存的硬上限,因为压缩、In-flight Request、对象元数据和其他结构还会占用额外内存。1
6.2 为什么需要 BufferPool
如果每次创建 Batch 都:
new byte[batch.size]高吞吐场景会持续创建和回收大 ByteBuffer,增加:
- 分配成本;
- GC 压力;
- 内存碎片;
- 尾延迟抖动。
BufferPool 复用固定或相近规模的 Buffer:
allocate ↓ProducerBatch ↓发送完成 / 失败终结 ↓deallocate ↓返回 Pool6.3 Buffer 为什么会耗尽
当持续满足:
Producer Append Rate λ>Sender Completion Rate μ待发送数据会增长。

常见原因:
- Broker 写入变慢;
- ACK 等待变长;
- 网络抖动;
- Metadata / Leader 不稳定;
- 请求不断重试;
- Broker 配额限制;
- 单条消息过大;
- Producer 瞬时流量峰值;
- Sender 线程 CPU 不足;
- DNS、TLS 或连接失败。
6.4 BufferPool 如何把背压传回业务线程
当无法复用 Batch 且 BufferPool 没有足够内存:
Application Thread ↓BufferPool.allocate() ↓等待其他 Batch 完成并释放 Buffer ↓在剩余 max.block.ms 内成功或超时失败这意味着 Producer 的背压最终会表现为:
send() latency ↑waiting threads ↑bufferpool-wait-time ↑TimeoutException / BufferExhaustedException6.5 增大 buffer.memory 能解决问题吗
它能吸收短时间突发:
短峰值 ↓Buffer 暂时积压 ↓峰值结束后 Sender 追平但无法解决持续:
λ > μ如果每秒新增积压 10 MiB,把 Buffer 从 32 MiB 改到 320 MiB,只是把失败时间从约 3 秒延后到约 30 秒,同时扩大进程内存和故障时未完成数据规模。
正确顺序是:
- 判断流量峰值还是持续过载;
- 定位 Sender、网络、Broker 或 ACK 瓶颈;
- 检查 Batch 是否过小;
- 检查 Partition 和 Broker 分布;
- 再决定 Buffer 是否需要扩展。
7. Sender:把 Partition Batch 变成 Broker Request
7.1 Sender 是 Producer Runtime 的推进器
业务线程只把 Record 放入 Accumulator。
Sender 必须不断完成:
Sender.run() ↓runOnce() ↓处理事务 / 幂等状态(存在时) ↓sendProducerData() ↓NetworkClient.poll()本篇只解释普通发送主线;PID、Sequence、事务状态机在第五篇深入。
7.2 Sender 核心循环
压缩后的逻辑:
while (running) { long pollTimeout = sendProducerData(now); client.poll(pollTimeout, now);}sendProducerData() 不只是简单调用 Socket,它需要先做调度:
检查过期 Batch ↓Accumulator.ready(metadata) ↓需要时触发 Metadata Refresh ↓选择 Ready Broker Nodes ↓Accumulator.drain(...) ↓构建 ProduceRequest ↓NetworkClient.send(...)
7.3 什么 Batch 才算 Ready
某个 Partition 的队首 Batch 可能因为以下条件变为 Ready:
- Batch 已满;
- 已等待至少
linger.ms; - Accumulator 内存紧张;
- Producer 正在 Flush;
- Producer 正在 Close;
- 事务状态要求发送;
- 其他触发立即 Drain 的状态。
但 Ready 还不等于可以发送。
Producer 还需要确认:
当前 Leader 已知?目标 Broker 连接 Ready?该 Node 的 In-flight Request 是否达到上限?Batch 是否还在 Retry Backoff?7.4 drain 为什么按 Broker Node 聚合
Accumulator 按 TopicPartition 存储:
orders-0orders-1payments-3但网络连接按 Broker Node 管理。
假设 Metadata:
orders-0 Leader → Broker 1orders-1 Leader → Broker 2payments-3 Leader → Broker 1Sender Drain 后需要得到:
Broker 1├── orders-0 Batch└── payments-3 Batch
Broker 2└── orders-1 Batch所以存在两个不同聚合层次:
ProducerBatch:TopicPartition 级ProduceRequest:Broker Node 级这是理解 Client Runtime 非常关键的一点。
7.5 ProduceRequest 可以包含多个 Batch
Kafka 协议以批量为中心,一次 ProduceRequest 可以为同一个 Broker 携带多个 Topic / Partition 的数据。
ProduceRequest Broker-1├── Topic orders│ ├── Partition 0 → Batch A│ └── Partition 4 → Batch B└── Topic audit └── Partition 2 → Batch C这进一步摊薄:
- Request Header;
- 网络系统调用;
- Broker 请求调度;
- TLS Record;
- 连接往返。
7.6 max.request.size 限制 Drain 结果
即使某个 Broker 有很多 Ready Batch,Sender 不能无限构建大 Request。
它需要在:
max.request.sizeBroker NodePartition OrderingAvailable Batches之间选择本轮发送集合。
未被 Drain 的 Batch 留待后续 Sender 循环。
7.7 Drain 不等于删除
Batch 从 Accumulator Drain 出来后,并不是立即释放内存。
它可能进入:
In-flight Request直到:
- 成功收到响应;
- 遇到不可恢复错误;
- 达到 Delivery Timeout;
- 重试状态终结;
- Producer 关闭并放弃。
可重试失败时,Batch 可能重新入队,因此不能在第一次 Network Send 后就销毁其状态。
8. NetworkClient:连接、In-flight Request 与响应完成
8.1 NetworkClient 承担什么职责
Sender 负责“哪些 Batch 应该发给谁”,NetworkClient 更接近通用 Kafka 协议网络层,负责:
Broker Connection StateDNS / Address ResolutionRequest Header / Correlation IDIn-flight RequestsSocket Read / WriteAuthentication / TLSResponse MatchingConnection FailureTimeout DetectionMetadata Request不要把 NetworkClient 理解成普通的:
socket.write(bytes);它维护一个持续运行的异步网络状态机。
8.2 Kafka Selector 与非阻塞 I/O
Java Client 通常使用单 I/O 线程配合非阻塞网络选择机制管理多个 Broker 连接:
Sender Thread ↓NetworkClient.poll() ↓Selector ├── Broker 1 Channel ├── Broker 2 Channel └── Broker 3 Channel这样一个 Producer 不需要为每个 Broker 创建一个业务发送线程。
8.3 In-flight Request 是什么
请求写入网络后、响应完成前:
ProduceRequest ↓In-flight ↓ProduceResponse / Timeout / DisconnectClient 需要记录:
- 目标 Node;
- Correlation ID;
- 创建时间;
- Request Timeout;
- 回调;
- 对应 Batch;
- 是否期待响应;
- 协议版本。
Broker 返回响应后,Client 根据 Correlation ID 找到对应请求。
8.4 max.in.flight.requests.per.connection
Kafka 4.3 默认:
max.in.flight.requests.per.connection = 5表示一个连接上最多允许多少个尚未收到 ACK 的请求。1
值更大可以提升流水线并行度:
Request 1 ───────── waitingRequest 2 ───────── waitingRequest 3 ───────── waiting但失败、重试和顺序状态会更复杂。
在关闭幂等生产且允许重试时,如果:
Batch A 先发 → 失败Batch B 后发 → 成功Batch A 重试 → 成功Broker Log 可能出现:
B → AKafka 4.3 官方文档明确提醒:当幂等关闭、重试开启且 In-flight 大于 1 时,可能发生重试导致的顺序变化;幂等开启时要求该值不超过 5,并在允许范围内保持顺序。1
PID、Sequence Number 和 Broker 去重窗口的具体机制放在第五篇。
8.5 NetworkClient 为什么也负责 Metadata
Metadata Request 本身也是 Kafka 协议请求,需要:
- 找到可连接 Broker;
- 建立连接;
- 发送请求;
- poll 响应;
- 更新 ProducerMetadata。
所以 Metadata 更新由 Sender / NetworkClient 推进,而等待 Metadata 的是业务线程。
这形成跨线程协作:
Application Thread ↓ wait metadataProducerMetadata ↑ updateSender / NetworkClient8.6 ProduceResponse 如何完成一个 Batch
Broker 响应通常包含每个 TopicPartition 的:
Error CodeBase OffsetLog Append TimeRecord Errors / Error Message(特定版本)Throttle TimeSender 根据响应把 Batch 分成三类。
成功
completeBatch ↓计算每条 RecordMetadata ↓完成 Future ↓执行 Callback ↓释放 Buffer可重试失败
例如某些:
Leader 变化临时网络断开Broker 暂时不可用请求超时可能执行:
标记 Metadata 需要刷新 ↓Batch 进入 Retry Backoff ↓重新入队 ↓后续 Sender Loop 再发不可恢复失败
例如:
消息过大序列化前已失败授权失败Topic 非法某些幂等 / 事务致命错误Delivery Timeout会完成 Future 为异常,并释放 Batch 资源。
8.7 Callback 在哪个线程执行
典型用法:
producer.send(record, (metadata, exception) -> { if (exception != null) { log.error("send failed", exception); return; } log.info("partition={}, offset={}", metadata.partition(), metadata.offset());});需要警惕:Callback 通常由 Producer I/O 推进线程执行。
因此不要:
producer.send(record, (metadata, exception) -> { slowDatabase.saveResult(metadata); // 慢 I/O});慢 Callback 可能阻塞 Sender 继续处理响应,进而使:
Network Poll 变慢Batch Completion 变慢Buffer 释放变慢send() Buffer Wait 增加更合理的方式是:
- Callback 只做轻量状态记录;
- 复杂处理转交专用 Executor;
- 严格限制队列与拒绝策略;
- 防止错误处理再次同步调用 Kafka 形成递归链路。
9. Four Timeout Boundaries:不要把所有超时混成一个参数
Producer 发送路径至少包含四类常被混淆的时间边界。

9.1 max.block.ms
Kafka 4.3 默认:
max.block.ms = 60000对普通 send() 来说,它主要限制:
等待 Metadata+等待 BufferPool 内存官方明确说明,用户 Serializer 或 Partitioner 中的阻塞时间不计入这一限制。1
所以:
send() 实际耗时可能 > max.block.ms如果是用户插件本身卡住。
9.2 linger.ms
Kafka 4.3 默认:
linger.ms = 5它是 Batch 未满时的批量等待上界之一,不是网络响应超时,也不是 Record 总交付期限。
9.3 request.timeout.ms
Kafka 4.3 默认:
request.timeout.ms = 30000它控制 Client 等待一次请求响应的最长时间;超时后可能:
- 进入可重试流程;
- 在重试耗尽或总期限到达后失败。
它不是从 send() 开始计算的全生命周期上限。
9.4 delivery.timeout.ms
Kafka 4.3 默认:
delivery.timeout.ms = 120000它覆盖从 send() 返回后,Record 在 Producer 内部等待、发送、等待 ACK 和可重试失败的总交付上界。1
概念关系:
Record Enter Producer ↓Accumulator Wait ↓ linger / schedulingNetwork Send ↓ request.timeoutRetry Backoff ↓More Attempts ↓Success or delivery.timeout.ms官方要求:
delivery.timeout.ms>=request.timeout.ms + linger.ms但生产设计不能只满足数学关系,还要根据:
- 业务 SLO;
- Broker 故障恢复时间;
- Leader 选举时间;
- 请求重试次数;
- 上游调用超时;
- 幂等要求;
- 未决 Record 数量;
共同规划。
9.5 为什么增加 Timeout 可能适得其反
假设上游 HTTP 请求超时 3 秒,但 Producer Delivery Timeout 是 2 分钟。
HTTP Client 3s 后放弃 ↓Producer 仍在后台重试 ↓最终消息可能成功 ↓调用方却认为业务失败并再次提交结果可能产生应用级重复。
所以端到端超时必须一起设计:
User Request TimeoutBusiness Transaction TimeoutProducer max.block.msProducer delivery.timeout.msBroker / Network TimeoutDownstream Idempotency10. Retry、Ordering 与 ACK:本篇需要掌握到什么深度
10.1 Retry 发生在哪里
Retry 不是业务线程重新调用 send(),而是 Sender 对可恢复 Batch 重新调度:
ProduceRequest ↓Retriable Error ↓Retry Backoff ↓Re-enqueue Batch ↓Refresh Metadata if Needed ↓Send AgainKafka 4.3 默认 retries 很大,官方建议主要使用 delivery.timeout.ms 控制总重试时间,而不是依赖手工设置有限重试次数。1
10.2 为什么应用级重试和 Client Retry 不同
Client Retry 可以保留同一 Producer Session、Batch 与幂等状态。
应用级重试通常是:
try { producer.send(record).get();} catch (Exception e) { producer.send(record); // 新的一次业务发送}这可能被 Producer 视为一条新的 Record,无法简单依赖 Client 内部去重。
官方 KafkaProducer 源码文档也提醒:幂等能力不能自动去重任意应用层重新发送,幂等保证有其 Producer Session 与协议边界。3
10.3 ACK 参数在执行链的位置
acks 影响的是:
ProduceRequest 何时被认为完成而不是:
Record 何时进入 Accumulator执行链:
RecordAccumulator ↓Sender ↓ProduceRequest(acks) ↓Broker Leader / Replicas ↓Response Completion ↓Future / Callback本篇只建立位置关系。
以下机制放在第五篇:
acks=0/1/all精确可靠性差异;- ISR;
min.insync.replicas;- Leader Crash;
- PID / Producer Epoch;
- Sequence Number;
- 重复检测;
- Transaction Coordinator;
- Exactly Once 边界。
10.4 幂等生产的现代默认值
Kafka 4.3 在无冲突配置时默认启用 enable.idempotence=true,并要求:
acks = allretries > 0max.in.flight.requests.per.connection <= 5如果显式配置冲突且又显式启用幂等,会触发配置错误;如果未显式启用,冲突配置可能导致幂等关闭。1
工程上不要为了“降低延迟”随意关闭幂等,必须先理解它对重试重复与顺序的影响。
11. 参数必须放回执行链理解
不要背一张“Kafka Producer 常用参数表”。
应该把参数绑定到组件和状态。
11.1 Metadata 与连接
| 参数 | 主要位置 | 解决的问题 |
|---|---|---|
bootstrap.servers | Bootstrap | 初次进入 Cluster |
metadata.max.age.ms | ProducerMetadata | 周期性强制刷新 |
metadata.max.idle.ms | Topic Metadata Cache | 空闲 Topic 缓存回收 |
metadata.recovery.strategy | Metadata Recovery | 已知 Broker 全不可用时是否 Rebootstrap |
client.dns.lookup | Address Resolution | Bootstrap DNS 使用方式 |
reconnect.backoff.ms | NetworkClient | 连接失败退避 |
11.2 Serialization 与消息规模
| 参数 | 主要位置 | 解决的问题 |
|---|---|---|
key.serializer | Application Thread | Key → bytes |
value.serializer | Application Thread | Value → bytes |
max.request.size | Size Check / Request Build | Request 与 Batch 大小边界 |
compression.type | RecordBatch Build | Batch 压缩 |
11.3 Partitioning
| 参数 | 主要位置 | 解决的问题 |
|---|---|---|
partitioner.class | Application Thread | 自定义 Partition 选择 |
partitioner.ignore.keys | Default Partition Logic | 是否忽略 Key |
partitioner.adaptive.partitioning.enable | Default Partitioner | 是否适应 Broker 表现 |
partitioner.availability.timeout.ms | Default Partitioner | 是否暂时避开不可用 Partition |
11.4 Batching 与 Memory
| 参数 | 主要位置 | 解决的问题 |
|---|---|---|
batch.size | ProducerBatch Allocation | 默认 Batch Buffer 大小 |
linger.ms | Sender Ready | 未满 Batch 等待上界 |
buffer.memory | BufferPool | 待发送数据内存预算 |
max.block.ms | Application Thread | Metadata + Buffer 等待上限 |
11.5 Network 与 Completion
| 参数 | 主要位置 | 解决的问题 |
|---|---|---|
max.in.flight.requests.per.connection | NetworkClient | 单连接并发未确认请求数 |
request.timeout.ms | In-flight Request | 单次请求响应等待 |
delivery.timeout.ms | Record / Batch Lifecycle | Record 总交付上界 |
retry.backoff.ms | Sender Retry | 重试退避 |
acks | ProduceRequest Completion | Broker 确认条件 |
enable.idempotence | Sender / Broker Protocol | Client 重试去重与顺序 |
11.6 参数的典型相互作用
batch.size × linger.ms × produce rate
低流量 + 大 batch.size + 小 linger→ Batch 仍可能较小
高流量 + 默认 linger→ 很快自然填满 Batchbuffer.memory × delivery latency
交付时间越长→ In-flight 与待发送数据占用越久→ Buffer 压力越大compression.type × batch fill ratio
Batch 更满→ 通常压缩上下文更丰富→ 压缩率可能更好max.in.flight × idempotence × retry
提高流水线并行度→ 必须同时考虑失败顺序和去重状态
12. Producer 调优:从目标和瓶颈出发
12.1 不存在一套通用“最优参数”
Producer 可能服务不同目标。
低延迟交易事件
关心:
P99 Delivery Latency可靠性顺序快速失败日志批量采集
关心:
ThroughputCompression RatioNetwork CostCPU Cost大事件或模型结果
关心:
Message SizeSerialization CPUMemory PressureBroker Limits跨地域发送
关心:
RTTIn-flight PipelineRetryBandwidthTLS CPU所以必须先确定:
吞吐目标平均 / 最大消息大小延迟 SLO失败容忍度顺序范围压缩收益Broker 容量12.2 吞吐优化的合理顺序
建议优先:
- 确认 Partition 数和 Broker 分布允许并行;
- 检查消息是否过小导致 Batch 填充差;
- 观察
batch-size-avg和records-per-request-avg; - 检查压缩是否降低网络瓶颈;
- 检查 Sender / Network Thread CPU;
- 检查 Broker Request Latency 与 Throttle;
- 检查 BufferPool Wait;
- 最后再扩大 Buffer 或 In-flight。
只增大 batch.size 但流量太低,不一定形成更大 Batch;可能只是预分配更多内存。
12.3 延迟优化不能只把 linger.ms=0
Kafka 4.0 将默认值改为 5 ms,原因之一是更有效的 Batch 可能减少请求数和系统负载,从而获得相近甚至更低的实际延迟。1
低延迟优化应同时观察:
record-queue-time-avgrequest-latency-avg / maxproduce-throttle-timebatch-size-avgrecord-retry-rateconnection setupBroker request queue如果真正瓶颈是 Broker 30 ms,减少 5 ms Linger 并不能解决根因。
12.4 内存优化不能只看 buffer.memory
Producer 实际内存还包括:
- ProducerBatch 元数据;
- Compression Buffer;
- Serializer 临时对象;
- In-flight Request;
- Network Buffers;
- Metrics;
- Callback 捕获对象;
- 业务 Record 在序列化前的对象;
- 自定义 Interceptor 状态。
所以需要同时观察:
JVM HeapAllocation RateGC PauseBuffer Available BytesWaiting ThreadsIn-flight Requests12.5 大消息不应该优先靠扩大限制解决
大消息会放大:
- Producer Buffer 占用;
- 压缩 CPU;
- Request 延迟;
- Broker Page Cache 污染;
- 副本网络;
- Consumer Fetch 内存;
- 重试成本;
- 单 Partition 尾延迟。
优先评估:
对象存储 + Kafka 引用事件拆分字段裁剪压缩格式Schema 优化而不是把所有相关限制从 1 MiB 一路改到 100 MiB。
13. Producer Metrics 与故障诊断
Kafka Java Client 提供内置 Metrics,并可通过 JMX 或自定义 MetricsReporter 输出。Kafka 4.3 官方监控文档说明 Java Client 使用内置 Kafka Metrics Registry,Rate 指标通常还有对应 -total 累计值。6
13.1 最值得关注的指标类别
Throughput
record-send-raterecord-send-totalbyte-raterequest-rateBatch
batch-size-avgbatch-size-maxrecords-per-request-avgcompression-rate-avgrecord-queue-time-avgBuffer
buffer-available-bytesbufferpool-wait-timewaiting-threadsNetwork
request-latency-avgrequest-latency-maxrequests-in-flightconnection-countconnection-creation-ratenetwork-io-rateError / Retry
record-error-raterecord-retry-raterecord-size-maxThrottle
produce-throttle-time-avgproduce-throttle-time-max具体指标名会随客户端版本和传感器层次变化,生产监控应从实际 4.3.x 客户端导出结果确认,而不是只复制旧版本列表。
13.2 发送变慢的根因树

可以先按执行链分段。
调用线程阶段
Serializer 慢?Interceptor 慢?Metadata Wait?BufferPool Wait?Accumulator 阶段
Batch 填充差?Record Queue Time 高?Buffer Available Bytes 下降?Network 阶段
连接失败?Request Timeout?Retry Rate 上升?In-flight 满?TLS CPU 高?Broker / ACK 阶段
Broker Request Latency?磁盘 / 副本延迟?Throttle?Leader 频繁变化?13.3 场景一:send() 本身变慢
优先看:
Metadata WaitBufferPool WaitSerializer / Interceptor CPUGC不要先看 Consumer Lag,因为消息可能还没离开 Producer。
13.4 场景二:Future 很久不完成,但 send() 很快
说明 Record 已顺利入队,问题更可能位于:
Batch SchedulingNetworkBrokerACKRetryDelivery Timeout13.5 场景三:Buffer Available 持续下降
排查:
生产速率是否持续大于交付速率?Broker 是否 Throttle?Request Latency 是否上升?Retry 是否上升?是否有连接故障?Callback 是否阻塞 Sender?13.6 场景四:吞吐低且 CPU 也低
可能是:
- Partition 太少;
- 业务发送速率本来就低;
- 每条后调用
flush(); - 每条 Future
.get(); - Batch 太小;
linger.ms过低;- 同步外部 Serializer / Interceptor;
- Broker 配额。
13.7 场景五:吞吐高但 P99 抖动
可能是:
- Batch / Compression CPU 抖动;
- GC;
- 大消息;
- Metadata Refresh;
- Leader 变化;
- Broker Throttle;
- Retry;
- TLS Handshake;
- Callback 执行过慢。
14. 可运行实验:观察 Producer 的 Batch、Partition 与背压
下面实验用于把 send → Accumulator → Sender → Broker 落到可观察行为。
环境定位:单节点、单副本、本地学习。该配置不代表生产建议。Docker 镜像和 Maven 依赖使用 Kafka 4.3.0,若本机镜像小版本不同,请保持 Client 与 Broker 在官方兼容范围内。
14.1 实验目录
kafka-producer-lab/├── docker-compose.yml├── pom.xml└── src/main/java/dev/jupiter/kafka/ └── ProducerLab.java14.2 Docker Compose
services: kafka: image: apache/kafka:4.3.0 container_name: kafka-ch03 hostname: kafka ports: - "9092:9092" environment: KAFKA_NODE_ID: 1 KAFKA_PROCESS_ROLES: broker,controller KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 volumes: - kafka-data:/var/lib/kafka/data
volumes: kafka-data:启动:
docker compose up -ddocker logs -f kafka-ch0314.3 创建 Topic
docker exec kafka-ch03 \ /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server localhost:9092 \ --create \ --topic producer-lab \ --partitions 6 \ --replication-factor 1查看:
docker exec kafka-ch03 \ /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server localhost:9092 \ --describe \ --topic producer-lab14.4 Maven 配置
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion>
<groupId>dev.jupiter</groupId> <artifactId>kafka-producer-lab</artifactId> <version>1.0.0</version>
<properties> <maven.compiler.release>21</maven.compiler.release> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> </properties>
<dependencies> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>4.3.0</version> </dependency> <dependency> <groupId>org.slf4j</groupId> <artifactId>slf4j-simple</artifactId> <version>2.0.17</version> </dependency> </dependencies>
<build> <plugins> <plugin> <groupId>org.codehaus.mojo</groupId> <artifactId>exec-maven-plugin</artifactId> <version>3.5.1</version> <configuration> <mainClass>dev.jupiter.kafka.ProducerLab</mainClass> </configuration> </plugin> </plugins> </build></project>14.5 Producer 实验代码
package dev.jupiter.kafka;
import org.apache.kafka.clients.producer.KafkaProducer;import org.apache.kafka.clients.producer.ProducerConfig;import org.apache.kafka.clients.producer.ProducerRecord;import org.apache.kafka.clients.producer.RecordMetadata;import org.apache.kafka.common.Metric;import org.apache.kafka.common.MetricName;import org.apache.kafka.common.serialization.StringSerializer;
import java.time.Duration;import java.util.Comparator;import java.util.Map;import java.util.Properties;import java.util.concurrent.CountDownLatch;import java.util.concurrent.TimeUnit;import java.util.concurrent.atomic.AtomicInteger;
public final class ProducerLab {
private static final String TOPIC = "producer-lab";
public static void main(String[] args) throws Exception { String mode = args.length == 0 ? "async" : args[0];
Properties props = new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ProducerConfig.CLIENT_ID_CONFIG, "producer-ch03-lab"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.ACKS_CONFIG, "all"); props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "zstd"); props.put(ProducerConfig.BATCH_SIZE_CONFIG, 64 * 1024); props.put(ProducerConfig.LINGER_MS_CONFIG, 20); props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 32L * 1024 * 1024); props.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 120_000); props.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, 30_000); props.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, 10_000);
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) { switch (mode) { case "sync" -> runSync(producer, 20_000); case "keyed" -> runKeyed(producer, 30); case "metrics" -> runAsync(producer, 100_000, true); default -> runAsync(producer, 100_000, false); }
producer.flush(); printMetrics(producer.metrics()); } }
private static void runAsync( KafkaProducer<String, String> producer, int count, boolean printProgress) throws InterruptedException {
CountDownLatch latch = new CountDownLatch(count); AtomicInteger failures = new AtomicInteger(); long start = System.nanoTime();
for (int i = 0; i < count; i++) { String key = "order-" + (i % 10_000); String value = "{\"seq\":" + i + ",\"type\":\"OrderCreated\"}";
producer.send(new ProducerRecord<>(TOPIC, key, value), (metadata, exception) -> { if (exception != null) { failures.incrementAndGet(); } latch.countDown(); });
if (printProgress && i % 10_000 == 0) { System.out.println("enqueued=" + i); } }
if (!latch.await(2, TimeUnit.MINUTES)) { throw new IllegalStateException("send completion timeout"); }
long elapsedMs = Duration.ofNanos(System.nanoTime() - start).toMillis(); System.out.printf("async count=%d elapsed=%dms failures=%d%n", count, elapsedMs, failures.get()); }
private static void runSync( KafkaProducer<String, String> producer, int count) throws Exception {
long start = System.nanoTime(); for (int i = 0; i < count; i++) { RecordMetadata metadata = producer.send( new ProducerRecord<>(TOPIC, "sync-key", "value-" + i) ).get();
if (i % 5_000 == 0) { System.out.printf("i=%d partition=%d offset=%d%n", i, metadata.partition(), metadata.offset()); } } long elapsedMs = Duration.ofNanos(System.nanoTime() - start).toMillis(); System.out.printf("sync count=%d elapsed=%dms%n", count, elapsedMs); }
private static void runKeyed( KafkaProducer<String, String> producer, int count) throws Exception {
for (int i = 0; i < count; i++) { String key = i % 2 == 0 ? "order-A" : "order-B"; RecordMetadata metadata = producer.send( new ProducerRecord<>(TOPIC, key, "event-" + i) ).get();
System.out.printf("key=%s partition=%d offset=%d%n", key, metadata.partition(), metadata.offset()); } }
private static void printMetrics(Map<MetricName, ? extends Metric> metrics) { metrics.entrySet().stream() .filter(entry -> { String name = entry.getKey().name(); return name.contains("batch-size") || name.contains("records-per-request") || name.contains("compression-rate") || name.contains("record-queue-time") || name.contains("request-latency") || name.contains("buffer-available") || name.contains("bufferpool-wait") || name.contains("record-retry") || name.contains("record-error"); }) .sorted(Comparator.comparing(entry -> entry.getKey().name())) .forEach(entry -> System.out.printf( "%s tags=%s value=%s%n", entry.getKey().name(), entry.getKey().tags(), entry.getValue().metricValue() )); }}代码使用 try-with-resources 关闭 Producer,Callback 保持轻量,并通过 Latch 等待异步发送完成。
14.6 实验一:Async 与每条 .get() 对比
异步:
mvn -q compile exec:java -Dexec.args="async"同步逐条等待:
mvn -q compile exec:java -Dexec.args="sync"观察:
总耗时batch-size-avgrecords-per-request-avgrequest-raterequest-latency预期:每条 .get() 会破坏流水线和 Batch 聚合能力。
注意:这不是说业务永远不能等待 Future,而是不能把高吞吐 Producer 误写成逐条同步 RPC。
14.7 实验二:验证同 Key 路由
mvn -q compile exec:java -Dexec.args="keyed"观察:
order-A 是否稳定进入同一 Partition?order-B 是否稳定进入同一 Partition?两个 Key 是否可能 Hash 到同一 Partition?两个不同 Key 映射到同一 Partition 是允许的,Hash 不是一对一。
14.8 实验三:观察 Batch Metrics
mvn -q compile exec:java -Dexec.args="metrics"分别修改:
batch.size=16384 / 65536 / 262144linger.ms=0 / 5 / 20 / 100compression.type=none / lz4 / zstd观察:
batch-size-avgrecords-per-request-avgcompression-rate-avgrecord-queue-time-avg总吞吐与 P99不要只比较单次运行。至少:
- 预热 JVM;
- 运行多轮;
- 控制消息大小;
- 控制 Topic Partition 数;
- 同时观察 Broker CPU 与网络。
14.9 实验四:制造 Buffer 背压
一种安全方式是本地临时降低 Broker 吞吐或暂停容器网络,而不是在生产环境操作。
可以:
- 把
buffer.memory改为较小值,例如 1 MiB; - 把
max.block.ms改为 2 秒; - 发送大量较大消息;
- 在发送过程中暂停 Broker:
docker pause kafka-ch03观察 Producer:
send() 是否开始阻塞?buffer-available-bytes 是否下降?bufferpool-wait-time 是否上升?最终异常是什么?恢复:
docker unpause kafka-ch03本实验只在本地学习环境执行。
14.10 实验五:观察 Metadata 错误
把 KAFKA_ADVERTISED_LISTENERS 暂时改为应用无法访问的地址,重启 Broker,再发送消息。
观察错误可能表现为:
Bootstrap 可以连接Metadata 可以获得但 Metadata 中返回的 Broker 地址不可达Producer 持续连接失败 / Timeout这能帮助理解:
Bootstrap 地址可达,不代表 Broker Advertised Address 对 Client 可达。
完成后恢复正确配置。
15. 源码阅读地图
15.1 第一层:公共 API
ProducerKafkaProducerProducerRecordRecordMetadataCallback关注:
send()返回 Future 的语义是什么?- 哪些异常可能同步抛出?
- 哪些异常通过 Future / Callback 返回?
flush()与close()如何影响未完成 Record?
15.2 第二层:KafkaProducer.doSend
建议围绕以下问题阅读:
Interceptor 在哪里执行?Metadata 等待在哪里发生?Serializer 何时调用?Partition 何时确定?Size Check 在哪里?Accumulator.append 输入是什么?什么条件下 sender.wakeup()?对应源码入口:
clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java15.3 第三层:Metadata
ProducerMetadataMetadataClusterPartitionInfo关注:
- Topic 是如何被加入待更新集合的?
- Application Thread 如何等待 Metadata Version 更新?
- Sender 如何知道需要发送 MetadataRequest?
- 错误如何使 Metadata 失效?
- Rebootstrap 如何触发?
15.4 第四层:Partitioning
KafkaProducer partition logicPartitionerBuiltInPartitionerCluster关注:
- 显式 Partition 在哪里校验?
- Key Hash 使用什么字节?
- Sticky Partition 何时切换?
- Adaptive Partitioning 依据什么可用状态?
- 自定义 Partitioner 和默认逻辑如何互斥?
15.5 第五层:Accumulator 与 Batch
RecordAccumulatorProducerBatchBufferPoolMemoryRecordsBuilder关注:
- 为什么数据结构是
TopicPartition → Deque? - Append 如何尝试复用 Batch?
- Buffer 在锁内还是锁外申请?为什么?
- Batch Full 如何判断?
- Future / Callback 如何挂到 Batch?
- Retry 如何重新入队?
- Complete 后何时释放 Buffer?
15.6 第六层:Sender
SenderTransactionManagerProduceRequest.Builder关注:
- Sender Loop 如何计算 Poll Timeout?
ready()输出哪些 Node?- 未知 Leader 如何触发 Metadata Update?
drain()如何按 Node 聚合?- Batch 何时标记 In-flight?
- Response 如何分类为 Success、Retry、Fail?
- Expired Batch 在哪里处理?
15.7 第七层:NetworkClient
NetworkClientInFlightRequestsClientRequestClientResponseSelectorKafkaChannel关注:
- Node Connection State 如何变化?
- Correlation ID 如何分配?
- In-flight 如何按 Node 管理?
- 请求超时在哪里检查?
- Disconnect 如何完成未决请求?
- Metadata Request 与 ProduceRequest 是否走同一网络基础设施?
15.8 推荐阅读顺序
不要从 KafkaProducer.java 第一行读到最后。
建议:
KafkaProducer API 文档 ↓ProducerRecord / RecordMetadata ↓KafkaProducer.send / doSend ↓ProducerMetadata ↓BuiltInPartitioner ↓RecordAccumulator.append ↓ProducerBatch ↓BufferPool ↓Sender.sendProducerData ↓NetworkClient.poll ↓ProduceResponse Completion15.9 固定源码阅读模板
每个方法都记录:
入口线程输入对象依赖状态锁与并发边界核心数据结构变化是否可能阻塞输出对象异步完成位置失败分类资源释放位置例如 RecordAccumulator.append:
线程:Application Thread输入:TopicPartition、Timestamp、Key/Value Bytes、Headers状态:Partition Deque、BufferPool、Batch Size动作:尝试追加现有 Batch,否则申请 Buffer 创建 Batch阻塞:可能等待 BufferPool输出:Future、Batch Created / Full 状态下游:Sender Ready / Drain释放:Batch 完成或失败终结后归还 Buffer16. 常见错误与反直觉结论
16.1 “send() 是异步的,所以绝不会影响业务线程”
错误。
Metadata Wait、Buffer Wait、Serializer、Interceptor、Partitioner 和同步异常都发生在调用线程。
16.2 “bootstrap.servers 必须列出全部 Broker”
错误。
它是 Bootstrap 入口,Client 会获取完整 Metadata。但生产应配置多个地址,避免初始节点故障。
16.3 “batch.size 就是 ProduceRequest 最大值”
错误。
Batch 是 TopicPartition 级;Request 是 Broker Node 级,可以包含多个 Batch。
16.4 “linger.ms=5 意味着每条消息都延迟 5 ms”
错误。
Batch 满、Flush、Memory Pressure 等条件可能提前 Ready;实际等待还可能因 Broker Backpressure 更长。
16.5 “Buffer 不够就无限增大 buffer.memory”
错误。
持续 λ > μ 时,任何有限 Buffer 最终都会耗尽。
16.6 “无 Key 用 Round Robin 最均匀,所以最好”
不完整。
逐条 Round Robin 可能破坏 Batch 填充。默认 Sticky / Adaptive 策略在批量效率和负载分布间做权衡。
16.7 “Callback 里可以随便写复杂业务”
错误。
慢 Callback 可能阻塞 Sender 的响应处理和 Buffer 释放。
16.8 “调用 Future.get() 才能保证可靠”
不完整。
.get() 让当前线程观察结果,但可靠性还取决于 ACK、ISR、幂等、Broker 配置和业务幂等;逐条 .get() 会显著削弱吞吐。
16.9 “设置 retries 很大就不会失败”
错误。
仍有不可恢复错误、Delivery Timeout、授权问题、消息过大、Producer 关闭和应用进程崩溃。
16.10 “Producer 成功回调等于下游业务处理成功”
错误。
它只表示 Produce 在 Kafka 定义的 ACK 边界内完成,不表示 Consumer 已消费,更不表示数据库事务或外部 API 已成功。
17. 面试题与场景题
17.1 基础问题
KafkaProducer.send()为什么是异步的?- Producer 为什么需要后台 Sender Thread?
Future<RecordMetadata>在什么时间完成?bootstrap.servers的真实作用是什么?- Producer 为什么需要 Metadata?
- Serializer 和 Partitioner 在哪个线程执行?
- 为什么 RecordAccumulator 按 TopicPartition 组织?
- ProducerBatch 和 ProduceRequest 的聚合单位分别是什么?
batch.size与max.request.size有什么区别?linger.ms何时生效?- BufferPool 为什么存在?
send()在什么情况下会阻塞?- Sender 如何把 Batch 按 Broker 聚合?
- NetworkClient 如何匹配 Request 和 Response?
- Callback 为什么不能执行慢 I/O?
17.2 深入问题
- 为什么 Metadata Wait 必须发生在 Partition 选择前?
- Key 的 Java
hashCode()是否决定 Kafka Partition? - 为什么无 Key 时 Sticky 策略可能比逐条 Round Robin 更快?
- 为什么一次 ProduceRequest 可以包含多个 TopicPartition Batch?
- Drain 之后为什么不能立即释放 Batch Buffer?
- BufferPool 耗尽如何体现 Broker Backpressure?
max.block.ms为什么不能限制 Serializer 卡死?delivery.timeout.ms与request.timeout.ms如何嵌套?- 开启 Retry、关闭 Idempotence 且 In-flight 大于 1 时为什么可能乱序?
- 应用重新调用
send()为什么不等价于 Client 内部 Retry? - Metadata 地址可达但 Broker 地址不可达时会发生什么?
- 为什么大
batch.size不一定产生大 Batch? - 为什么增大 Buffer 不能解决持续过载?
- Producer 发送延迟高时如何按执行链排查?
17.3 场景题一:高峰期 send() 频繁超时
已知:
buffer-available-bytes 持续下降bufferpool-wait-time 上升Broker request latency 上升回答应包括:
Broker / Network 交付速度下降 ↓Batch 长时间占用 Buffer ↓BufferPool 耗尽 ↓Application Thread 等待 ↓max.block.ms 超时不能只回答“把 Buffer 调大”。
17.4 场景题二:发送吞吐很低但 Broker 很空闲
代码:
for (...) { producer.send(record).get();}问题是逐条同步等待破坏 Pipeline 和 Batch。
优化方向:
- 异步 Send;
- Callback / 批量等待;
- 有界并发;
- 最后 Flush;
- 同时保持错误可观测。
17.5 场景题三:同一订单事件出现顺序错乱
需要检查:
是否使用相同 Key?Key Serializer 是否一致?Topic 是否扩过 Partition?是否存在应用级重发?是否关闭 Idempotence?max.in.flight 与 Retry 如何配置?是否跨多个 Producer / Topic?Consumer 是否并发改变业务顺序?17.6 场景题四:第一次发送总是慢
检查:
- Metadata Bootstrap;
- DNS;
- TLS / SASL;
- Topic 自动创建;
- Connection Setup;
- Serializer 类加载;
- Producer 是否频繁重建。
常见解决方案是长期复用 Producer,并在应用就绪阶段完成必要预热,而不是每个请求创建 Producer。
17.7 场景题五:Producer CPU 很高
检查:
- Compression Codec 与 Level;
- Serializer;
- Interceptor;
- 小 Batch 导致高 Request Rate;
- TLS;
- Metrics Recording Level;
- 大量 Callback;
- 频繁创建 Producer;
- GC 与对象分配。
18. 本篇验收清单
完成本篇后,你应当能够脱稿完成以下任务。
18.1 画图验收
画出:
Application Threads ↓KafkaProducer.send ↓Metadata / Serialize / Partition ↓RecordAccumulator ↓ProducerBatch Deque
================ Thread Boundary ================
Sender ↓Ready / Drain ↓ProduceRequest by Broker Node ↓NetworkClient / In-flight ↓Broker Leader ↓Response ↓Future / Callback / Buffer Release18.2 概念验收
可以明确区分:
ProducerRecordSerialized RecordProducerBatchRecordBatchProduceRequestIn-flight RequestRecordMetadata18.3 参数验收
可以把以下参数放回执行链,而不是单独背定义:
bootstrap.serversbatch.sizelinger.msbuffer.memorymax.block.msmax.request.sizecompression.typerequest.timeout.msdelivery.timeout.msmax.in.flight.requests.per.connectionacksenable.idempotence18.4 故障验收
面对发送超时,先问:
卡在 Application Thread?Metadata?Serializer?BufferPool?Accumulator?Sender Scheduling?Network Connection?Broker Processing?ACK Waiting?Retry / Delivery Timeout?18.5 最终口述验收
你应该能不看文章,用 30 分钟讲清:
一条
ProducerRecord如何从 Java 业务线程,经过 Metadata、Serialization、Partition Selection 和 RecordAccumulator,进入 ProducerBatch;Sender 又如何按 Broker Node Drain Batch、构建 ProduceRequest、通过 NetworkClient 发送,并在响应后完成 Future 与释放 Buffer。
如果只能说出:
Producer 序列化消息,然后发到 Broker说明仍未掌握 Producer Runtime。
19. 全文收束:一条 Record 的完整旅程
19.1 Application Thread
ProducerRecord ↓ProducerInterceptor.onSend ↓waitOnMetadata(topic) ↓Serialize Key / Value ↓Partition Selection ↓Record Size Check ↓RecordAccumulator.append ↓Future<RecordMetadata>19.2 Accumulator
TopicPartition ↓Locate Deque Tail Batch ↓Append Existing Batch或BufferPool Allocate ↓Create ProducerBatch ↓Append and Wake Sender19.3 Sender
Batch Full / Linger / Flush / Memory Pressure ↓Accumulator.ready ↓Resolve Partition Leader ↓Drain Batches by Broker Node ↓Build ProduceRequest ↓NetworkClient.send19.4 Network 与 Completion
In-flight Request ↓Broker ProduceResponse ↓Success / Retriable / Fatal ↓Complete or Re-enqueue Batch ↓Future / Callback ↓BufferPool deallocateKafka Producer 的高吞吐并不来自业务线程“更快地调用 Socket”。
它来自一套清晰的异步 Pipeline:
应用线程负责把业务 Record 转换为按 Partition 聚集的待发送 Batch;Sender 线程负责把这些 Batch 按 Broker 路由、形成网络流水线,并用响应驱动 Future、Retry 与内存回收。
从工程角度看,Producer 不是一个简单 API,而是一个运行在应用进程中的小型消息发送系统:
Metadata Cache+Partition Router+Batching Buffer+Memory Pool+Async Scheduler+Network State Machine+Retry / Timeout State+Metrics理解这套 Runtime 后,Producer 调优就不再是盲目修改参数。
你会知道:
batch.size属于 Batch Allocation;linger.ms属于 Ready Scheduling;buffer.memory属于 Client Backpressure;max.block.ms属于 Application Thread;max.in.flight属于 Network Pipeline;request.timeout.ms属于单次 Request;delivery.timeout.ms属于 Record Lifecycle;acks属于 Broker Completion Boundary。
20. 下一篇预告:Consumer Group、Coordinator 与 Rebalance
Producer 的核心难题是:
如何把大量小 Record 高效、可靠地路由到 Partition Leader?
Consumer 的核心难题则完全不同:
多个动态加入和离开的 Consumer,如何协调 Partition Ownership、消费进度和故障恢复?
下一篇将从:
consumer.poll(Duration.ofSeconds(1));继续追踪:
- Consumer 为什么采用 Pull;
- Fetch Position、Committed Offset、High Watermark 和 Last Stable Offset 如何区分;
- Group Coordinator 保存什么状态;
- FindCoordinator、JoinGroup、SyncGroup 和 Heartbeat 如何协作;
- Consumer 为什么会 Rebalance;
- Classic Protocol 为什么有全局同步屏障;
- Kafka 4.x 新 Consumer Rebalance Protocol 为什么转向 Broker-side Assignment 与 Incremental Coordination;
- Share Groups 与传统 Consumer Groups 的模型差异。
下一篇:
《Kafka Consumer Group:一个分布式成员协调系统的设计与实现》
参考资料
- Apache Kafka 4.3 Producer Configs
- Apache Kafka 4.3 KafkaProducer API
- Apache Kafka KafkaProducer Source
- Apache Kafka 4.3 ProducerInterceptor API
- Apache Kafka 4.3 Design
- Apache Kafka 4.3 Monitoring
- Apache Kafka RecordAccumulator Source
- Apache Kafka ProducerBatch Source
- Apache Kafka Sender Source
- Apache Kafka NetworkClient Source
- Apache Kafka ProducerMetadata Source
- Apache Kafka BufferPool Source