12671 字
63 分钟
Kafka Producer 完整执行链:一条 Record 从 send() 到 Partition Leader 的旅程

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.mslinger.msrequest.timeout.msdelivery.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

send 返回不等于消息已经发出

Kafka 官方 API 将 send() 定义为异步操作:Record 进入待发送 Buffer 后,调用可以返回,从而允许 Producer 将多条 Record 合并成更高效的批次。2

但要立刻补一个重要边界:

异步发送不等于 send() 永远不会阻塞。

调用线程至少可能在以下位置停留:

  1. Topic Metadata 尚不可用,需要等待 Metadata 更新;
  2. BufferPool 内存不足,需要等待 Sender 释放 Buffer;
  3. 用户自定义 Interceptor、Serializer 或 Partitioner 本身执行很慢;
  4. Record 过大或配置非法,可能在调用线程直接失败;
  5. 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
等待 Metadata
Serializer
Partition 选择
Record Size 校验
RecordAccumulator Append
返回 Future

Sender I/O Thread#

执行:

检查 Batch 是否 Ready
处理过期 Batch
刷新 Metadata
按 Broker Node Drain Batch
构建 ProduceRequest
管理连接
发送请求
poll 网络事件
处理响应或错误
完成 Callback / Future

Kafka Producer 双线程协作模型

线程边界带来三个关键能力。

第一,批处理#

多个业务线程调用 send() 时,Record 可以先聚集到按 Partition 组织的 Batch 中,再由一个 Sender 统一发送。

第二,调用与网络解耦#

业务线程不需要为每条 Record 执行一次完整 Broker 往返。

第三,显式背压边界#

如果 Sender 长期发送不过去,BufferPool 会被消耗,最终使新的业务线程在 send() 内等待甚至失败。

这说明 Producer 并没有“消灭背压”,只是把背压从同步 RPC 响应时间转换为:

Accumulator Queue Depth
BufferPool Available Memory
Batch Age
Request Latency
Retry / Timeout

1.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>

KafkaProducer.send 主执行链

这里最重要的认知是:

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
);

它可能包含:

Topic
Optional Partition
Optional Timestamp
Optional Key
Value
Headers

此时 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.size
Broker message.max.bytes
Topic max.message.bytes
Replica fetch size
Consumer 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

Bootstrap 与 Metadata Refresh

3.2 Producer 为什么可以直连 Leader#

Kafka 不要求所有 Produce 流量经过一个中央 Router。

Producer 获取 Metadata 后,可以建立映射:

Topic: order-events
P0 Leader → Broker 1
P1 Leader → Broker 3
P2 Leader → Broker 2

然后:

orders-P1 Batch
直接发送 Broker 3

Kafka 官方设计明确说明,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 State

Producer 最关注:

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-events

Producer 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 Partition

Partitioner 不是一个“负载均衡小工具”,而是业务语义到分布式 Log 的映射函数。

4.2 默认选择顺序#

现代 Kafka 默认逻辑可以概括为:

1. Record 显式指定 Partition
→ 使用显式 Partition
2. 未指定 Partition,但存在 Key
→ 根据 Key Hash 选择 Partition
3. 既没有 Partition,也没有 Key
→ 使用 Sticky / Adaptive 选择策略形成更有效 Batch

Partition 选择决策树

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 = orderId
OrderCreated(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 → P0
R2 → P1
R3 → P2
R4 → 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>>

RecordAccumulator 按 TopicPartition 组织批次

例如:

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 Batch

Request 总上限更多受 max.request.size 约束。

5.5 大于 batch.size 的单条 Record 怎么办#

官方配置说明指出,Producer 不会尝试把大于 batch.size 的 Record 与其他 Record 按普通方式合批,但这不意味着它一定被 batch.size 拒绝。Client 可以为单个大 Record 分配足够容纳它的 Buffer,只要没有超过更高层的消息和请求大小限制。1

所以:

batch.size
单条消息绝对最大值

真正需要联合检查:

estimated serialized size
max.request.size
Broker / Topic message limits
Buffer availability

5.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 1
Record 2
Record 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 C

Broker 对 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
返回 Pool

6.3 Buffer 为什么会耗尽#

当持续满足:

Producer Append Rate λ
>
Sender Completion Rate μ

待发送数据会增长。

BufferPool 与 Producer 背压

常见原因:

  • 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 / BufferExhaustedException

6.5 增大 buffer.memory 能解决问题吗#

它能吸收短时间突发:

短峰值
Buffer 暂时积压
峰值结束后 Sender 追平

但无法解决持续:

λ > μ

如果每秒新增积压 10 MiB,把 Buffer 从 32 MiB 改到 320 MiB,只是把失败时间从约 3 秒延后到约 30 秒,同时扩大进程内存和故障时未完成数据规模。

正确顺序是:

  1. 判断流量峰值还是持续过载;
  2. 定位 Sender、网络、Broker 或 ACK 瓶颈;
  3. 检查 Batch 是否过小;
  4. 检查 Partition 和 Broker 分布;
  5. 再决定 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(...)

Sender Ready、Drain 与 Request 构建

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-0
orders-1
payments-3

但网络连接按 Broker Node 管理。

假设 Metadata:

orders-0 Leader → Broker 1
orders-1 Leader → Broker 2
payments-3 Leader → Broker 1

Sender 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.size
Broker Node
Partition Ordering
Available 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 State
DNS / Address Resolution
Request Header / Correlation ID
In-flight Requests
Socket Read / Write
Authentication / TLS
Response Matching
Connection Failure
Timeout Detection
Metadata 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 / Disconnect

Client 需要记录:

  • 目标 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 ───────── waiting
Request 2 ───────── waiting
Request 3 ───────── waiting

但失败、重试和顺序状态会更复杂。

在关闭幂等生产且允许重试时,如果:

Batch A 先发 → 失败
Batch B 后发 → 成功
Batch A 重试 → 成功

Broker Log 可能出现:

B → A

Kafka 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 metadata
ProducerMetadata
↑ update
Sender / NetworkClient

8.6 ProduceResponse 如何完成一个 Batch#

Broker 响应通常包含每个 TopicPartition 的:

Error Code
Base Offset
Log Append Time
Record Errors / Error Message(特定版本)
Throttle Time

Sender 根据响应把 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 发送路径至少包含四类常被混淆的时间边界。

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 / scheduling
Network Send
↓ request.timeout
Retry 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 Timeout
Business Transaction Timeout
Producer max.block.ms
Producer delivery.timeout.ms
Broker / Network Timeout
Downstream Idempotency

10. Retry、Ordering 与 ACK:本篇需要掌握到什么深度#

10.1 Retry 发生在哪里#

Retry 不是业务线程重新调用 send(),而是 Sender 对可恢复 Batch 重新调度:

ProduceRequest
Retriable Error
Retry Backoff
Re-enqueue Batch
Refresh Metadata if Needed
Send Again

Kafka 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 = all
retries > 0
max.in.flight.requests.per.connection <= 5

如果显式配置冲突且又显式启用幂等,会触发配置错误;如果未显式启用,冲突配置可能导致幂等关闭。1

工程上不要为了“降低延迟”随意关闭幂等,必须先理解它对重试重复与顺序的影响。


11. 参数必须放回执行链理解#

不要背一张“Kafka Producer 常用参数表”。

应该把参数绑定到组件和状态。

11.1 Metadata 与连接#

参数主要位置解决的问题
bootstrap.serversBootstrap初次进入 Cluster
metadata.max.age.msProducerMetadata周期性强制刷新
metadata.max.idle.msTopic Metadata Cache空闲 Topic 缓存回收
metadata.recovery.strategyMetadata Recovery已知 Broker 全不可用时是否 Rebootstrap
client.dns.lookupAddress ResolutionBootstrap DNS 使用方式
reconnect.backoff.msNetworkClient连接失败退避

11.2 Serialization 与消息规模#

参数主要位置解决的问题
key.serializerApplication ThreadKey → bytes
value.serializerApplication ThreadValue → bytes
max.request.sizeSize Check / Request BuildRequest 与 Batch 大小边界
compression.typeRecordBatch BuildBatch 压缩

11.3 Partitioning#

参数主要位置解决的问题
partitioner.classApplication Thread自定义 Partition 选择
partitioner.ignore.keysDefault Partition Logic是否忽略 Key
partitioner.adaptive.partitioning.enableDefault Partitioner是否适应 Broker 表现
partitioner.availability.timeout.msDefault Partitioner是否暂时避开不可用 Partition

11.4 Batching 与 Memory#

参数主要位置解决的问题
batch.sizeProducerBatch Allocation默认 Batch Buffer 大小
linger.msSender Ready未满 Batch 等待上界
buffer.memoryBufferPool待发送数据内存预算
max.block.msApplication ThreadMetadata + Buffer 等待上限

11.5 Network 与 Completion#

参数主要位置解决的问题
max.in.flight.requests.per.connectionNetworkClient单连接并发未确认请求数
request.timeout.msIn-flight Request单次请求响应等待
delivery.timeout.msRecord / Batch LifecycleRecord 总交付上界
retry.backoff.msSender Retry重试退避
acksProduceRequest CompletionBroker 确认条件
enable.idempotenceSender / Broker ProtocolClient 重试去重与顺序

11.6 参数的典型相互作用#

batch.size × linger.ms × produce rate#

低流量 + 大 batch.size + 小 linger
→ Batch 仍可能较小
高流量 + 默认 linger
→ 很快自然填满 Batch

buffer.memory × delivery latency#

交付时间越长
→ In-flight 与待发送数据占用越久
→ Buffer 压力越大

compression.type × batch fill ratio#

Batch 更满
→ 通常压缩上下文更丰富
→ 压缩率可能更好

max.in.flight × idempotence × retry#

提高流水线并行度
→ 必须同时考虑失败顺序和去重状态

Producer 参数调优 Trade-off


12. Producer 调优:从目标和瓶颈出发#

12.1 不存在一套通用“最优参数”#

Producer 可能服务不同目标。

低延迟交易事件#

关心:

P99 Delivery Latency
可靠性
顺序
快速失败

日志批量采集#

关心:

Throughput
Compression Ratio
Network Cost
CPU Cost

大事件或模型结果#

关心:

Message Size
Serialization CPU
Memory Pressure
Broker Limits

跨地域发送#

关心:

RTT
In-flight Pipeline
Retry
Bandwidth
TLS CPU

所以必须先确定:

吞吐目标
平均 / 最大消息大小
延迟 SLO
失败容忍度
顺序范围
压缩收益
Broker 容量

12.2 吞吐优化的合理顺序#

建议优先:

  1. 确认 Partition 数和 Broker 分布允许并行;
  2. 检查消息是否过小导致 Batch 填充差;
  3. 观察 batch-size-avgrecords-per-request-avg
  4. 检查压缩是否降低网络瓶颈;
  5. 检查 Sender / Network Thread CPU;
  6. 检查 Broker Request Latency 与 Throttle;
  7. 检查 BufferPool Wait;
  8. 最后再扩大 Buffer 或 In-flight。

只增大 batch.size 但流量太低,不一定形成更大 Batch;可能只是预分配更多内存。

12.3 延迟优化不能只把 linger.ms=0#

Kafka 4.0 将默认值改为 5 ms,原因之一是更有效的 Batch 可能减少请求数和系统负载,从而获得相近甚至更低的实际延迟。1

低延迟优化应同时观察:

record-queue-time-avg
request-latency-avg / max
produce-throttle-time
batch-size-avg
record-retry-rate
connection setup
Broker 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 Heap
Allocation Rate
GC Pause
Buffer Available Bytes
Waiting Threads
In-flight Requests

12.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-rate
record-send-total
byte-rate
request-rate

Batch#

batch-size-avg
batch-size-max
records-per-request-avg
compression-rate-avg
record-queue-time-avg

Buffer#

buffer-available-bytes
bufferpool-wait-time
waiting-threads

Network#

request-latency-avg
request-latency-max
requests-in-flight
connection-count
connection-creation-rate
network-io-rate

Error / Retry#

record-error-rate
record-retry-rate
record-size-max

Throttle#

produce-throttle-time-avg
produce-throttle-time-max

具体指标名会随客户端版本和传感器层次变化,生产监控应从实际 4.3.x 客户端导出结果确认,而不是只复制旧版本列表。

13.2 发送变慢的根因树#

Producer 发送延迟诊断树

可以先按执行链分段。

调用线程阶段#

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 Wait
BufferPool Wait
Serializer / Interceptor CPU
GC

不要先看 Consumer Lag,因为消息可能还没离开 Producer。

13.4 场景二:Future 很久不完成,但 send() 很快#

说明 Record 已顺利入队,问题更可能位于:

Batch Scheduling
Network
Broker
ACK
Retry
Delivery Timeout

13.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.java

14.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:

启动:

Terminal window
docker compose up -d
docker logs -f kafka-ch03

14.3 创建 Topic#

Terminal window
docker exec kafka-ch03 \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--create \
--topic producer-lab \
--partitions 6 \
--replication-factor 1

查看:

Terminal window
docker exec kafka-ch03 \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--describe \
--topic producer-lab

14.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() 对比#

异步:

Terminal window
mvn -q compile exec:java -Dexec.args="async"

同步逐条等待:

Terminal window
mvn -q compile exec:java -Dexec.args="sync"

观察:

总耗时
batch-size-avg
records-per-request-avg
request-rate
request-latency

预期:每条 .get() 会破坏流水线和 Batch 聚合能力。

注意:这不是说业务永远不能等待 Future,而是不能把高吞吐 Producer 误写成逐条同步 RPC。

14.7 实验二:验证同 Key 路由#

Terminal window
mvn -q compile exec:java -Dexec.args="keyed"

观察:

order-A 是否稳定进入同一 Partition?
order-B 是否稳定进入同一 Partition?
两个 Key 是否可能 Hash 到同一 Partition?

两个不同 Key 映射到同一 Partition 是允许的,Hash 不是一对一。

14.8 实验三:观察 Batch Metrics#

Terminal window
mvn -q compile exec:java -Dexec.args="metrics"

分别修改:

batch.size=16384 / 65536 / 262144
linger.ms=0 / 5 / 20 / 100
compression.type=none / lz4 / zstd

观察:

batch-size-avg
records-per-request-avg
compression-rate-avg
record-queue-time-avg
总吞吐与 P99

不要只比较单次运行。至少:

  • 预热 JVM;
  • 运行多轮;
  • 控制消息大小;
  • 控制 Topic Partition 数;
  • 同时观察 Broker CPU 与网络。

14.9 实验四:制造 Buffer 背压#

一种安全方式是本地临时降低 Broker 吞吐或暂停容器网络,而不是在生产环境操作。

可以:

  1. buffer.memory 改为较小值,例如 1 MiB;
  2. max.block.ms 改为 2 秒;
  3. 发送大量较大消息;
  4. 在发送过程中暂停 Broker:
Terminal window
docker pause kafka-ch03

观察 Producer:

send() 是否开始阻塞?
buffer-available-bytes 是否下降?
bufferpool-wait-time 是否上升?
最终异常是什么?

恢复:

Terminal window
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#

Producer
KafkaProducer
ProducerRecord
RecordMetadata
Callback

关注:

  1. send() 返回 Future 的语义是什么?
  2. 哪些异常可能同步抛出?
  3. 哪些异常通过 Future / Callback 返回?
  4. 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.java

15.3 第三层:Metadata#

ProducerMetadata
Metadata
Cluster
PartitionInfo

关注:

  1. Topic 是如何被加入待更新集合的?
  2. Application Thread 如何等待 Metadata Version 更新?
  3. Sender 如何知道需要发送 MetadataRequest?
  4. 错误如何使 Metadata 失效?
  5. Rebootstrap 如何触发?

15.4 第四层:Partitioning#

KafkaProducer partition logic
Partitioner
BuiltInPartitioner
Cluster

关注:

  1. 显式 Partition 在哪里校验?
  2. Key Hash 使用什么字节?
  3. Sticky Partition 何时切换?
  4. Adaptive Partitioning 依据什么可用状态?
  5. 自定义 Partitioner 和默认逻辑如何互斥?

15.5 第五层:Accumulator 与 Batch#

RecordAccumulator
ProducerBatch
BufferPool
MemoryRecordsBuilder

关注:

  1. 为什么数据结构是 TopicPartition → Deque
  2. Append 如何尝试复用 Batch?
  3. Buffer 在锁内还是锁外申请?为什么?
  4. Batch Full 如何判断?
  5. Future / Callback 如何挂到 Batch?
  6. Retry 如何重新入队?
  7. Complete 后何时释放 Buffer?

15.6 第六层:Sender#

Sender
TransactionManager
ProduceRequest.Builder

关注:

  1. Sender Loop 如何计算 Poll Timeout?
  2. ready() 输出哪些 Node?
  3. 未知 Leader 如何触发 Metadata Update?
  4. drain() 如何按 Node 聚合?
  5. Batch 何时标记 In-flight?
  6. Response 如何分类为 Success、Retry、Fail?
  7. Expired Batch 在哪里处理?

15.7 第七层:NetworkClient#

NetworkClient
InFlightRequests
ClientRequest
ClientResponse
Selector
KafkaChannel

关注:

  1. Node Connection State 如何变化?
  2. Correlation ID 如何分配?
  3. In-flight 如何按 Node 管理?
  4. 请求超时在哪里检查?
  5. Disconnect 如何完成未决请求?
  6. 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 Completion

15.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 完成或失败终结后归还 Buffer

16. 常见错误与反直觉结论#

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 基础问题#

  1. KafkaProducer.send() 为什么是异步的?
  2. Producer 为什么需要后台 Sender Thread?
  3. Future<RecordMetadata> 在什么时间完成?
  4. bootstrap.servers 的真实作用是什么?
  5. Producer 为什么需要 Metadata?
  6. Serializer 和 Partitioner 在哪个线程执行?
  7. 为什么 RecordAccumulator 按 TopicPartition 组织?
  8. ProducerBatch 和 ProduceRequest 的聚合单位分别是什么?
  9. batch.sizemax.request.size 有什么区别?
  10. linger.ms 何时生效?
  11. BufferPool 为什么存在?
  12. send() 在什么情况下会阻塞?
  13. Sender 如何把 Batch 按 Broker 聚合?
  14. NetworkClient 如何匹配 Request 和 Response?
  15. Callback 为什么不能执行慢 I/O?

17.2 深入问题#

  1. 为什么 Metadata Wait 必须发生在 Partition 选择前?
  2. Key 的 Java hashCode() 是否决定 Kafka Partition?
  3. 为什么无 Key 时 Sticky 策略可能比逐条 Round Robin 更快?
  4. 为什么一次 ProduceRequest 可以包含多个 TopicPartition Batch?
  5. Drain 之后为什么不能立即释放 Batch Buffer?
  6. BufferPool 耗尽如何体现 Broker Backpressure?
  7. max.block.ms 为什么不能限制 Serializer 卡死?
  8. delivery.timeout.msrequest.timeout.ms 如何嵌套?
  9. 开启 Retry、关闭 Idempotence 且 In-flight 大于 1 时为什么可能乱序?
  10. 应用重新调用 send() 为什么不等价于 Client 内部 Retry?
  11. Metadata 地址可达但 Broker 地址不可达时会发生什么?
  12. 为什么大 batch.size 不一定产生大 Batch?
  13. 为什么增大 Buffer 不能解决持续过载?
  14. 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 Release

18.2 概念验收#

可以明确区分:

ProducerRecord
Serialized Record
ProducerBatch
RecordBatch
ProduceRequest
In-flight Request
RecordMetadata

18.3 参数验收#

可以把以下参数放回执行链,而不是单独背定义:

bootstrap.servers
batch.size
linger.ms
buffer.memory
max.block.ms
max.request.size
compression.type
request.timeout.ms
delivery.timeout.ms
max.in.flight.requests.per.connection
acks
enable.idempotence

18.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 Sender

19.3 Sender#

Batch Full / Linger / Flush / Memory Pressure
Accumulator.ready
Resolve Partition Leader
Drain Batches by Broker Node
Build ProduceRequest
NetworkClient.send

19.4 Network 与 Completion#

In-flight Request
Broker ProduceResponse
Success / Retriable / Fatal
Complete or Re-enqueue Batch
Future / Callback
BufferPool deallocate

Kafka 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:一个分布式成员协调系统的设计与实现》


参考资料#

  1. Apache Kafka 4.3 Producer Configs
  2. Apache Kafka 4.3 KafkaProducer API
  3. Apache Kafka KafkaProducer Source
  4. Apache Kafka 4.3 ProducerInterceptor API
  5. Apache Kafka 4.3 Design
  6. Apache Kafka 4.3 Monitoring
  7. Apache Kafka RecordAccumulator Source
  8. Apache Kafka ProducerBatch Source
  9. Apache Kafka Sender Source
  10. Apache Kafka NetworkClient Source
  11. Apache Kafka ProducerMetadata Source
  12. Apache Kafka BufferPool Source
Kafka Producer 完整执行链:一条 Record 从 send() 到 Partition Leader 的旅程
https://jupiter-ws.cn/posts/backend/kafka/03_kafka_producer_execution_chain/
作者
Jupiter
发布于
2026-07-17
许可协议
CC BY-NC-SA 4.0