有些Kafka的源码分析文章,上来就贴一堆类名和方法签名,看完除了记住了几个名词,脑子里还是浆糊。我一开始读KafkaProducer的时候也是这个状态,后来踩了几个线上问题回头看,才慢慢把整条链路串起来。这篇文章我不打算把每个方法源码逐行贴一遍,那没意义,我更想带你把一条消息从send()出去,到 Broker 确认返回,这条完整路径上的关键设计、核心源码逻辑、以及参数为什么这么调,一次讲透。搞懂这一条链,Kafka 很多调优和报错问题你都能自己推到答案。
这套机制解决的核心问题其实就一个:网络 IO 太慢,不能让业务线程等。所以 Kafka Producer 把发送拆成了两段——业务线程只管把消息丢进内存缓冲,一个独立的后台线程负责批量打包、网络发送、接收响应。这个异步模型是理解全部源码的地基。适合谁看?准备深入 Kafka 的 Java 开发者,或者被线上“消息延迟高”“发送超时”这类问题折磨过、想从源头搞明白原因的同学。
1. 先把发送链路画在脑子里:Kafka Producer 的整体设计
1.1 这套异步架构到底解决了什么问题
先想一个场景:你的业务线程调用一次producer.send(),如果这条消息要立刻建立 TCP 连接、立刻写 socket、然后阻塞等 Broker 返回 ack,整套流程下来一次要多少时间?快则几毫秒,慢则几十上百毫秒。而业务接口往往一次请求要发好几条消息,如果全是同步阻塞,接口延迟直接爆炸。
所以 Kafka 选择了异步缓冲模型。业务线程调send()的时候,实际只做了几件轻量的事:序列化、分区计算、把消息追加到一个内存批次(RecordBatch)里,然后赶紧返回。数据攒在内存里,由一个叫 Sender 的后台线程统一调度,攒够一批或者到了时间阈值再批量发出去。大量小消息合并成一次网络请求,吞吐量能提升几个数量级。
这个设计思路在源码里的体现就是:KafkaProducer负责入口逻辑,RecordAccumulator管消息攒批缓冲,Sender线程管网络发送。每一条消息从生产到发送完成,必然依次经过这三层。读源码如果脑子里没有这条主线,很容易迷失在KafkaProducer那几百行注释和一堆内部类里。
1.2 源码阅读的地图:核心类各管哪一段
我在读源码之前,建议你先做一件事:把下面这个对应关系抄下来或者记住,后面读源码的时候随时回来对照。
| 核心类 | 职责 | 对应链路阶段 |
|---|---|---|
KafkaProducer | 对外的生产者入口,负责序列化、分区、调用累加器 | 消息进入管道 |
Partitioner | 决定消息进哪个分区 | 分区选择 |
RecordAccumulator | 内存中攒批,维护每个分区的批次队列 | 攒批缓冲 |
BufferPool | 管理内存块(ByteBuffer)的分配和回收 | 内存管理 |
Sender | 后台发送线程,批量取出批次并发送 | 发送调度 |
NetworkClient | 底层网络通信,管理连接和请求/响应 | 网络IO |
Metadata | 维护集群元数据,主题分区信息,分区 leader 所在节点 | 元数据驱动 |
一句话概括消息的流向:业务线程往累加器里放,Sender线程从累加器里取,NetworkClient负责真正把数据发到Broker。
有了这个地图,接下来每一步我们都能定位到具体类和方法,就不会在源码的海洋里漂着。我下面会按消息的实际流动顺序,从send()入口开始,一步步往底层走,同时把每一步的关键参数和设计原因讲清楚。
2. 从 send() 到 RecordAccumulator:消息进管道前的每一步
2.1 追踪 doSend():拦截器、序列化器、分区器各干了什么
send()方法本身是接口方法,真正干活的是内部私有的doSend()。我用简化代码把主流程标出来,你对照源码看会更清晰:
private Future<RecordMetadata> doSend(ProducerRecord<K, V> record, Callback callback) { TopicPartition tp = null; try { // 1. 先拦截器过一遍,允许拦截器修改消息或者做附加处理 ProducerRecord<K, V> interceptedRecord = this.interceptors.onSend(record); // 2. 如果这个主题的元数据还没有,先等元数据 ClusterAndWaitTime clusterAndWaitTime = waitOnMetadata(interceptedRecord.topic(), interceptedRecord.partition(), maxBlockTimeMs); // 3. 序列化 key 和 value byte[] serializedKey = keySerializer.serialize(interceptedRecord.topic(), interceptedRecord.headers(), interceptedRecord.key()); byte[] serializedValue = valueSerializer.serialize(interceptedRecord.topic(), interceptedRecord.headers(), interceptedRecord.value()); // 4. 计算分区 int partition = partition(interceptedRecord, serializedKey, serializedValue, cluster, tp); tp = new TopicPartition(interceptedRecord.topic(), partition); // 5. 追加到累加器,返回 future RecordAccumulator.RecordAppendResult result = accumulator.append(tp, interceptedRecord.timestamp(), serializedKey, serializedValue, interceptedRecord.headers(), interceptedRecord.key(), interceptedRecord.value(), interceptCallback, nowMaxBlockTimeMs); // 6. 如果批次满了或者新建了批次,唤醒 Sender 线程发送 if (result.batchIsFull || result.newBatchCreated) { sender.wakeup(); } return result.future; } }注意一个细节:序列化和分区都排在拦截器后面,所以拦截器里看到的是原始对象,而之后的操作基于序列化结果。拦截器适合做什么?链路追踪、消息脱敏、附加审计信息,这些都行。但拦截器本身是同步执行的,千万别在拦截器里做耗时操作,比如查数据库或者远程调用,那会把整个发送线程堵死。这个坑我踩过,生产环境消息吞吐暴跌,最后定位到是拦截器里调了个外部 HTTP 接口做风控,延迟平均加了 20 毫秒。
2.2 元数据拉取:为什么要先等 metadata
第二个关键点在第 2 步:waitOnMetadata。为什么发送消息要先拿元数据?因为 Producer 得知道这条消息该发给哪个 Broker——你这主题有 8 个分区,分区 0 的 leader 在节点 A,分区 1 的 leader 在节点 B,不查元数据怎么知道网络请求往哪里发?
Metadata内部维护了一个 Cluster 对象,里面有所有 topic 的 partition 信息、leader 节点信息。这个对象不是一次性拉全的,而是 Producer 启动时先拉一次全量,之后通过后台更新机制或者发送失败时的强制刷新来保持最新。如果你的主题第一次发消息,元数据里没有这个 topic,Producer 会主动向任意一个 Broker 请求元数据,并且带着max.block.ms的超时等待。这就是为什么有时候你刚建好 Topic 立刻生产,第一条消息会稍微慢一点,因为要等元数据刷新。
这里有个面试常问的坑:如果这个 Topic 在服务器上根本不存在,且allow.auto.create.topics默认是 true,那 Broker 那边会自动创建主题并返回元数据。但如果你在生产环境禁用了自动创建,而你的代码又写错了 Topic 名,send()不会立刻报错,而是会阻塞等待直到max.block.ms超时,然后抛TimeoutException: Topic not present in metadata after 60000 ms。我见过不少新手第一次碰到这个报错一脸懵,就是没理解 Producer 发消息前要过元数据这一关。
2.3 分区计算的逻辑:有 key 和无 key 差别很大
分区计算的核心逻辑在DefaultPartitioner里。理解它之前,先记住一条主线:Producer 拍板消息进哪个分区,然后往该分区 leader 所在的 Broker 发数据。
分区的规则很简单,分两种情况:
- 指定了 partition:那没得说,直接用指定的分区,源码里连分区器都不会调用。
- 没指定 partition 但是有 key:对 key 的字节数组做 murmur2 哈希,然后对分区数取模。同一个 key 永远会进同一个分区,这是 Kafka 保证分区有序的关键。
- 没指定分区也没 key:走粘性分区(Sticky Partitioning)策略。随机选一个分区,然后在这个分区上攒一批,攒满了再换下一个。这个策略是 Kafka 2.4 引入的,目的就是解决老版本一个问题——每条消息随机选分区导致批次太小,压缩率和批量发送效率都上不去。粘性分区让一批消息尽量落在同一个分区,显著提升批次填充率和吞吐量。
这段逻辑虽然不复杂,但很能说明 Kafka 的一个设计原则:能少做决策就少做决策,把工作集中在能批量化的地方。你如果自己做分区策略,也建议遵循这个思路——先保证 key 的哈希均匀性,再考虑怎么提高批次命中率。
3. RecordAccumulator:高吞吐背后的内存池设计
3.1 为什么要搞一个 BufferPool,而不是直接 new ByteBuffer
很多人在看RecordAccumulator之前,觉得它就是一个队列集合:每个分区对应一个队列,队列里存批次,发一条消息就往队列尾巴上追加。这个理解方向没错,但它漏掉了一个核心设计——内存管理。Kafka 在 Java 堆内维护了一个专门的BufferPool来管理发送缓冲的 ByteBuffer,这背后是为了避免一个老生常谈的问题:GC。
试想一下,如果你的业务每秒发送几十万条消息,意味着每秒要创建几十万个字节数组,用完又丢掉。JVM 堆内对象越多,GC 压力越大,最直观的表现就是 Kafka Producer 频繁 Full GC,然后消费者那侧开始报“连接被重置”。Kafka 的做法是池化内存:申请固定大小的内存块(默认 16KB),用双端队列管理空闲块,用完就归还到池里复用,尽量不让 JVM 频繁分配新对象。
BufferPool的核心代码逻辑不算复杂,重点是这两个方法:
allocate(int size, long maxTimeToBlockMs):从池里取一个 ByteBuffer。如果没有空闲块,就阻塞等待其他线程释放。deallocate(ByteBuffer bb):把用完的 ByteBuffer 归还到池里。
整个 Pool 由一把ReentrantLock保护,用Condition管理等待线程。如果batch.size配置小于等于默认池的单个块大小,那么整个批次就是一把梭,直接申请一块池内存;如果消息太大,超过单块大小,池里兜不住,就退化为直接ByteBuffer.allocate分配堆外内存。这个细节对应一个常见问题:batch.size设得太大,内存池的意义就打折扣了,GC 压力反而上来。
3.2 RecordBatch 的组装过程:一次追加,三次尝试
消息进入累加器后,会尝试追加到一个已经存在的RecordBatch里。RecordBatch这个类是整个发送链路最核心的数据结构之一,它内部维护了一个MemoryRecordsBuilder,真正的消息数据是以一种压缩二进制格式存储的,不是纯 Java 对象。把消息序列化成字节之后,通过DefaultRecordBatch的格式写入底层 ByteBuffer。
看accumulator.append()的实现,会发现它用了三次尝试的循环:
// 第一次尝试 synchronized (tp) { // 如果该分区已经有 batch,试图追加 if (batch.tryAppend(...)) { return result; } } // 如果批次满了或者没有批次,尝试申请新内存 // 1. 尝试从空闲的 batch 池复用 // 2. 从 BufferPool 分配新内存 // 然后 new RecordBatch 并追加为什么循环三次?核心是因为多线程并发。多个业务线程同时往同一个分区的 batch 追加,当前 batch 可能在竞争下满了,但满的那一刻没人知道,等锁释放后必须重新检查;同理,重新申请的内存可能在分配的过程中被其他线程用掉了,又要重新走一遍。这个 try 循环的设计告诉我们一个经验:多线程场景下的资源追加,不能用“判断+操作”两个独立步骤,必须在一个临界区内完成。
tryAppend里还有个细节:如果这个 batch 已经有压缩过的数据了,新消息还能不能追加?答案是能。它会把已有数据解压(如果设置了max.in.flight等),把新消息推进去,再重新压缩。所以压缩不是攒满一批才开始压,而是一边追加一边压缩。这也是为什么compression.type设置会显著影响 CPU 消耗和吞吐量——你设置在 Producer 端压缩,数据在内存里就是压缩状态,网络传输量和 Broker 磁盘存储量都能降下来。
3.3 batch.size、linger.ms、buffer.memory 三个参数是怎么协作的
这是面试出镜率最高的三个参数,也是调优绕不开的三件套。我直接说结论,然后解释源码里怎么体现的。
batch.size:单个批次的最大字节数,默认 16KB。它决定一条消息什么时候触发“批次满了,赶紧发”。linger.ms:批次在内存里最多等多久,默认 0 毫秒(即来一条发一条,但实际不会那么极端,因为线程调度)。它决定一条消息最长能憋多久。buffer.memory:Producer 总共用来缓冲待发送消息的内存大小,默认 32MB。它决定背压的阈值。
发送的触发条件是:batch 满了 或者 linger 超时了 或者 producer 要关闭了。在源码里,RecordAccumulator.append()返回的RecordAppendResult带了一个batchIsFull标志,如果追加后这个批次确实满了,主线程就调用sender.wakeup(),把 Sender 线程从poll阻塞中唤醒,赶紧把这个批次取走发送。这就是linger.ms=0为什么也能攒批的原因——只要多个线程在同一毫秒内往同一个 batch 追加,batch 可能很快就满了。
我做过一组测试,给大家一个直观数据:单条消息 200 字节,batch.size16KB,linger.ms=0,吞吐量大约在 8 万条/秒;linger.ms调到 5ms,吞吐量能到 30 万条/秒左右,CPU 占用反而更低。原因很简单,一个批次里塞进去的消息更多,网络请求次数更少,分摊到每条消息的压缩、传输、协议开销都降了。但如果你的场景对延迟极敏感,比如就要毫秒级投递,那linger.ms就别乱调了,维持默认或者设 1ms 就行。
buffer.memory相对独立。它管的是最坏情况下积压多少数据。如果业务突发写入量超过网络发送能力,内存里会堆积,等到全满之后,再调用send()就会阻塞,阻塞超过max.block.ms(默认 60 秒)就抛异常。这是 Kafka 的保护机制,防止内存无限膨胀把进程打垮。调大这个值能扛更大的峰值,但代价是延迟增加和 GC 压力上升,得结合你实际的堆积场景来定。
4. Sender 线程与 NetworkClient:数据真正走出内存
4.1 Sender 线程:循环里的三次选择
RecordAccumulator只是把消息缓冲好了,真正决定“什么时候发、发给谁、发哪些批次”的是Sender线程。它启动之后就一直在一个 while 循环里运行,核心逻辑是runOnce()。这个方法里的关键步骤,我帮大家简化成“三次选择”:
- 第一次选择:从元数据封装中筛出哪些节点是“可写”的(
ready()方法) - 第二次选择:每个可写节点,选出哪些分区的批次要发给它(
drain()方法) - 第三次选择:对每个批次,决定在什么超时时间内发送(
sendProducerData()方法)
ready()方法判断一个 broker 是否 ready 的条件很多,核心是:分区 leader 在这个 broker 上、有数据要发、且满足发送条件。发送条件包括muted状态、inFlightRequests的容量限制等。
drain()方法的作用是尽量减少网络请求次数。它按节点聚合所有要发给这个节点的批次,攒成一个列表。这里有一个重要逻辑:如果多个分区 leader 都在同一个 broker,那它们的数据很可能合并到一个请求里,这也是 Kafka 能高效处理几百个分区的关键。
4.2 in-flight 请求与 max.in.flight.requests.per.connection
NetworkClient是真正干网络活的地方。它管理了所有到 Broker 的连接、请求的发送和响应的接收。其中有个关键结构:InFlightRequests,它维护了每个节点当前“在途”的请求数量。
max.in.flight.requests.per.connection这个参数默认是 5,意思是每个 Broker 连接上最多允许 5 个还没收到响应的请求。如果设成 1,配合retries > 0,可以严格保证分区内消息的顺序。为什么?因为同一时刻只有一个请求在飞,这个请求失败了,下一个请求还没发出去,重试的话不会出现后面的请求先到了 Broker 这种乱序情况。代价是吞吐量下降,因为单连接单请求在飞,链路空闲的时间增多了。
源码里判断一个节点是否 ready 时,会检查 in-flight 数是否达到上限:
if (inFlightRequests.isFull(node)) { // 写缓冲压力大,暂不发送 return false; }这里有个实战经验:如果你的业务允许少量乱序,建议不要为了“保险”把max.in.flight.requests.per.connection设为 1。Kafka 2.x 之后有个更细粒度的控制方式:开启enable.idempotence=true(幂等生产者),配合retries=Integer.MAX_VALUE,在保证顺序的同时还能维持较高的吞吐。幂等生产者靠序列号和 PID 在 Broker 端去重,解决了网络重试导致的重复消息问题。这个参数组合已经成了现代 Kafka 生产者的默认推荐配置。
4.3 请求发送的细节:sendProducerData() 与超时管理
sendProducerData()做的事情是:遍历ready()选出的节点,对每个节点调用drain()取出该发的批次列表,然后通过NetworkClient.send()发送ProduceRequest。这里面有几个值得注意的细节:
- 一个节点只有一个 TCP 连接,Kafka 的协议是纯异步的,同一个连接上可以同时存在多个未完成的请求。
- 每个
ProduceRequest里可以包含多个分区的数据,这就是批量发送的网络层体现。 - 在 send 之前,
ProducerBatch上会记录一个createdMs时间戳和produceFuture回调,等响应回来之后通过这些元数据触发回调。
再往下走,NetworkClient的底层就是 Java NIO 的Selector,它把 SocketChannel 注册到 selector 上,然后poll()监听可读可写事件。Sender 线程的 while 循环每次都会调client.poll(...),这个方法会处理:
- 发送之前需要进行的连接建立(包括等待元数据时发起的连接)
- 已经写出去、等待响应的请求的在途管理
- 收到响应的数据处理(比如更新元数据、完成 batch 回调)
如果你打开源码看看poll()的逻辑,会发现它的核心就是遍历 selector 的事件桶,然后分别调用handleCompletedSends()、handleCompletedReceives()、handleDisconnections()、handleConnections()。每一个方法对应网络生命周期的一个事件。顺着这个结构读,源码就不乱了。
5. 响应、重试与回调:一条消息的收尾工作
5.1 ProduceResponse 回来之后发生了什么
请求发出去之后,NetworkClient收到 Broker 的响应,会封装成一个ClientResponse。在 Sender 线程的handleCompletedReceives()中,最终会走到completeBatch()方法。这一步做的事基本就是两件:
- 解析
ProduceResponse,拿到每个分区对应的错误码和 offset - 把结果交回给消息生产者侧的回调机制
这里要特别注意:回调(Callback)并不是在业务线程执行,而是在 Sender 线程里执行的。也就是说,producer.send()传入的 callback 一旦被调用,你永远不应该在 callback 里做耗时操作或者发阻塞请求(比如再调一次同步send()),因为这会直接阻塞 Sender 线程,导致整个 Producer 的所有发送操作变慢。
这个坑我遇到过一次非常典型的情况:某服务用Callback里实时把发送结果写进日志表,日志表的写延迟偶尔飙到几百毫秒,结果整个 Kafka 生产吞吐量被拖垮,引发了连环故障。后来改成 callback 里只更新内存计数器,另起线程批量落盘,问题立刻消失。
5.2 重试机制:什么时候重试,会不会乱序
发送失败后,Kafka 会尝试重试。重试的逻辑在Sender里,核心逻辑是:如果 batch 发送失败,且retries > 0,并且这条消息没有超过delivery.timeout.ms的总时限,就把它重新放回累加器,等待下次调度重新发送。
retries参数控制重试次数,默认值是 2147483647(配合幂等开启时)。retries=0表示不重试,网络抖动一下消息就丢了。delivery.timeout.ms默认 120000,即消息必须在 2 分钟内送达,超了就放弃。这个“总时限”是累加linger.ms、重试等待时间、请求超时时间得出来的,你调优的时候别只盯着request.timeout.ms和retries,它们是并列关系。
一个常见的乱序场景:retries > 0且max.in.flight.requests.per.connection > 1。假设你连续发了两条消息到同一个分区,请求 1 失败但请求 2 成功,重试请求 1 成功后,Broker 端实际写入的顺序是 2 然后 1,这就乱序了。这就是为什么很多老手会把 inflight 设为 1 来保序,或者更推荐的方式是开启幂等后让它自动做序列化顺序保证。
5.3 回调执行与异常处理:Future 和 Callback 各走各的路
每个send()调用会返回一个Future<RecordMetadata>,这就是doSend()里第 5 步 append 的返回值之一。注意这个 Future 不是 JDK 原生的FutureTask,而是FutureRecordMetadata,它的get()方法会阻塞等待结果,或者直接抛出发送异常。
如果你调用future.get(),异常会在这里抛出。如果你传了 Callback,异常会在 Sender 线程内通过拦截器链的onError传递。这两条路径互不干扰但结果相等——同一个发送动作,异常要么在 Future 上抛,要么在 Callback 里收到,不会两个都触发,这点源码里用了一个try-finally结构的保证。
还有一个很好的设计点:如果消息最终发送成功,RecordMetadata里会包含 topic、partition、offset、timestamp 这些信息,方便业务侧做审计或者对账。如果失败,异常类型包括RecordTooLargeException(消息单条超过max.request.size)、TimeoutException、KafkaException等,你可以根据异常类型决定是否需要重试,而不是一梭子打死。
6. 生产环境实战:常见报错与调优排查记录
6.1 高频报错一:Topic not present in metadata 与 cluster authorization failed
这两个报错是 Kafka 开发者最常碰到的两个首轮报错。先看症状:
org.apache.kafka.common.errors.TimeoutException: Topic test-topic not present in metadata after 60000 ms.原因大概率是:allow.auto.create.topics被设为 false(或集群侧没开自动建 Topic),而你的 Topic 名错误,或者消息发到了还没创建成功的分区。处理方式:先去服务端用命令行创建 Topic 并确认名称,别凭印象写。这个报错因为要等max.block.ms超时,经常被误报为“网络超时”,实际上跟网络没啥关系。
另一个高频报错:
org.apache.kafka.common.errors.ClusterAuthorizationException: Cluster authorization failed.这个报错跟 ACL 有关,通常是用户权限不足。有些团队基于不同环境用同一套代码跑,测试环境的用户有 admin 权限,一道生产环境就报这个错,然后疯狂排查网络。实际上就是生产环境的 ACL 策略没给这个用户分配生产权限。办法是找 Kafka 管理员确认该用户名和 Topic 的权限关系,重点检查Describe、Write、Create这几项权限是否齐全。
6.2 线上消息延迟高,怎么一步步定位
线上“消息延迟高”是排查非常难缠的问题,因为它可能出在 Producer、Broker、Consumer 任意一环。如果 Producer 侧发送延迟高,建议按下面这个顺序排查:
- 先看发送队列堆积情况:观察
buffer.memory的占用率,如果长期高位,说明生产速度大于网络发送能力,要么是网络带宽瓶颈,要么是压缩配置没开。 - 再看 Sender 线程忙不忙:如果 Sender 线程在等待 IO,可能是 Broker 侧的
socket.send.buffer.bytes太小或者分区数分配不均导致的请求排队。 - 看单批次大小:如果生产消息单条非常大(比如超过 100KB),
batch.size默认 16KB 会导致每条消息单独占一个批次,批量效应失效,网络请求翻倍,延迟自然高。这时候应该调大batch.size或者调大max.request.size,注意两者都要调。 - 最后看客户端 GC:如果内存池设计和 GC 没配合好,频繁 GC 也会导致发送线程卡顿。观察 GC 日志,如果 Young GC 频率离谱,考虑调大
batch.size减少对象分配次数。
我遇到过一个真实案例:某业务消息平均 300KB,batch.size保持默认 16KB,结果每条消息独立成批,单个请求只发一条消息,网络开销巨大,线上 TPS 从 2 万掉到 2000。后来把batch.size调到 1MB、max.request.size调到 1MB、启用 lz4 压缩,吞吐量直接恢复。这个案例说明:参数调优不能离开具体消息大小谈,同样一套配置对不同大小的消息效果天差地别。
6.3 高并发场景的调优组合建议
结合我自己的生产经验,给几个可以直接抄作业的组合场景:
| 场景 | 推荐配置 | 理由 |
|---|---|---|
| 高吞吐优先,允许少量延迟 | acks=all,linger.ms=5~10,batch.size=32KB~64KB,compression.type=lz4 | 批次更大,压缩省带宽,吞吐最高 |
| 低延迟优先 | acks=1,linger.ms=0,batch.size=16KB | 消息尽量不积压,来一条发一条 |
| 必须严格有序 | enable.idempotence=true,max.in.flight.requests.per.connection=1 | 幂等 + 单飞行请求保证分区内顺序 |
| 峰值冲击大 | buffer.memory=64MB~128MB,max.block.ms=5000 | 大缓冲扛峰值,但及时暴露阻塞 |
还有一个被好多人忽略的配置:max.request.size。它控制单个请求的最大字节数。如果一条消息本身就很大,超过这个限制,会直接抛RecordTooLargeException,而且不会重试。调它时记得把batch.size也调大,否则可能出现“批次装不下单条消息”的隐性问题。
另外,生产环境一定要开监控。Kafka Producer 的 JMX 指标里,record-queue-time-avg、record-send-rate、request-latency-avg这几个是最核心的,配合 grafana 面板做基线报警,延迟突然升高能第一时间发现。不要等业务方反馈“消息慢了”才去查,到那时候通常已经积累了很久了。
我个人做 Kafka Producer 调优这么多年,最大的体会是:不要试图记住每个参数默认值,而要理解一条消息从发出到确认的每一步消耗在哪儿。延迟高,就沿着链路找哪一步耗时最长;吞吐低,就看批次有没有攒满;乱序,就问自己重试和 inflight 设置是不是互相打架。把源码链路吃透了,这些问题你自己就能推导出答案,比背任何参数表都管用。如果这篇文章能帮你把 Producer 的源码链路串起来,我的目的就达到了。