Producer 原理与核心配置
producer.send() 一行代码背后,消息经历了拦截器、序列化、分区路由、内存攒批、网络发送五个阶段。理解这条链路,你才能解释"为什么消息没有立刻发出去"、"为什么内存爆了"、"batch.size 和 linger.ms 到底怎么配"。
1. 发送流程全景
Producer 内部是双线程模型:主线程负责把消息放进缓冲区,Sender 线程负责真正的网络 IO。
主线程 (你的业务线程)
send(record)
|
v
[拦截器 Interceptors] onSend() 加工消息 (打标、埋点)
|
v
[序列化器 Serializer] key/value 对象 -> byte[]
|
v
[分区器 Partitioner] 决定去哪个分区 (有 key 按哈希)
|
v
[RecordAccumulator 缓冲区] 按 "topic-分区" 维度攒批
| buffer.memory=32MB batch1[msg,msg,..] batch2[..]
| |
- - | - - - - - - - - - - - - - - - - | - - - - - -
v v
Sender 线程 (IO 线程) 批次满(batch.size) 或到时(linger.ms)
把就绪批次组装成请求 ──> 网络发送到分区 leader 所在 broker
等待 acks 响应 ──> 成功: 执行回调 / 失败: 重试或抛错关键点:
send()返回时消息只是进了本地缓冲区,并没有发出去。- 一个批次(batch)只属于一个分区,发往同一 broker 的多个批次会合并成一个请求。
- Sender 线程异步工作,主线程几乎无阻塞——除非缓冲区满。
2. 攒批三参数
| 参数 | 默认值 | 含义 |
|---|---|---|
batch.size | 16384 (16KB) | 单个批次的字节上限,攒满立即可发 |
linger.ms | 0 | 批次未满时最多等多久,到时也发 |
buffer.memory | 33554432 (32MB) | 整个 Producer 缓冲区总大小 |
发送时机 = 批次攒满 batch.size 或 等待达到 linger.ms,先到先发。
linger.ms=0(默认):消息尽快发出,延迟最低,但批次很小、吞吐差。linger.ms=5~100:故意等一等换取更大的批次,吞吐大幅提升,延迟增加几毫秒——高吞吐场景必调。- 缓冲区满时
send()会阻塞,最长max.block.ms(默认 60 秒),超时抛异常。
linger.ms=0: [m][m][m][m] 每条几乎单独一个请求, 请求数多
linger.ms=20: [m m m m m m m][...] 攒成大批次, 请求数骤减, 吞吐up3. 压缩
compression.type 对整个批次压缩,批次越大压缩比越高(这也是要攒批的原因之一):
| 算法 | 压缩比 | CPU 开销 | 建议 |
|---|---|---|---|
| none | - | 无 | 默认 |
| lz4 | 中 | 低 | 通用首选 |
| snappy | 中 | 低 | 与 lz4 接近 |
| zstd | 高 | 中 | 带宽敏感首选(2.1+) |
| gzip | 高 | 高 | 不推荐新项目 |
4. Java 客户端实战
依赖 org.apache.kafka:kafka-clients:3.7.1。
import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.Properties;
import java.util.concurrent.Future;
public class OrderProducer {
public static void main(String[] args) throws Exception {
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
// 攒批与压缩
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 32768); // 32KB
props.put(ProducerConfig.LINGER_MS_CONFIG, 20); // 等 20ms
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 64 * 1024 * 1024L);
Producer<String, String> producer = new KafkaProducer<>(props);
ProducerRecord<String, String> record =
new ProducerRecord<>("orders", "user-1", "{\"orderId\":1001}");
// 方式1: 发后即忘 —— 最快, 可能丢消息且无感知
producer.send(record);
// 方式2: 同步发送 —— 逐条确认, 吞吐最低
RecordMetadata meta = producer.send(record).get();
System.out.printf("sync -> partition=%d offset=%d%n",
meta.partition(), meta.offset());
// 方式3: 异步 + 回调 —— 生产推荐, 高吞吐且能感知失败
producer.send(record, (metadata, exception) -> {
if (exception != null) {
// 记日志/告警/落补偿表, 不要吞掉
exception.printStackTrace();
} else {
System.out.printf("async -> partition=%d offset=%d%n",
metadata.partition(), metadata.offset());
}
});
producer.flush(); // 强制清空缓冲区
producer.close(); // 会先 flush 再关闭
}
}⚠️Producer 是线程安全的, 全局一个就够
KafkaProducer 线程安全且很重——内部有缓冲区和 IO 线程。正确用法是整个应用共享一个实例。每次发送 new 一个 Producer 是最常见的性能反模式,会导致连接风暴且完全无法攒批。另外,进程退出前必须调用 close(),否则缓冲区里还没发出去的消息会直接丢失。
5. Python 对照(confluent-kafka)
from confluent_kafka import Producer
producer = Producer({
"bootstrap.servers": "localhost:9092",
"batch.size": 32768,
"linger.ms": 20,
"compression.type": "lz4",
})
def on_delivery(err, msg):
if err is not None:
print(f"delivery failed: {err}")
else:
print(f"delivered to {msg.topic()}[{msg.partition()}]@{msg.offset()}")
for i in range(10):
producer.produce(
topic="orders",
key=f"user-{i % 3}",
value=f'{{"orderId": {1000 + i}}}',
on_delivery=on_delivery,
)
producer.poll(0) # 触发已完成发送的回调
producer.flush() # 退出前必须 flushproduce() 同样只是入队;poll(0) 驱动回调执行;flush() 阻塞直到全部送达。
6. 其他重要参数速查
| 参数 | 默认 | 说明 |
|---|---|---|
max.request.size | 1MB | 单个请求最大字节数,超大消息会被直接拒绝 |
request.timeout.ms | 30000 | 等 broker 响应的超时 |
delivery.timeout.ms | 120000 | send() 到确认成功/失败的总时限(含重试) |
max.block.ms | 60000 | 缓冲区满或元数据不可用时 send() 最长阻塞时间 |
client.id | 空 | 客户端标识,方便服务端日志与配额定位 |
💡吞吐调优三板斧
高吞吐场景先调这三个:linger.ms=20~100、batch.size=64KB~128KB、compression.type=lz4/zstd。通常能把吞吐提升 3-10 倍,代价只是几十毫秒的发送延迟。可靠性相关的 acks、retries、幂等在下一章专门讲。
小结
- Producer 双线程模型:主线程写缓冲区,Sender 线程网络发送;
send()返回不等于已发送 - 链路:拦截器 → 序列化 → 分区器 → RecordAccumulator 攒批 → Sender
- 攒批由
batch.size与linger.ms共同控制,先到先发;压缩作用于整个批次 - 三种发送:发后即忘(慎用)、同步(低吞吐)、异步回调(生产推荐)
- 全局共享一个 Producer 实例,退出前 close/flush
- 下一章解决"怎么保证消息不丢不重" →
🎯练习
- 写一个循环发送 10 万条消息的程序,分别用
linger.ms=0和linger.ms=50+ lz4 压缩跑一遍,对比总耗时。 - 把
buffer.memory调小到 1MB 并快速发送大消息,观察 send() 阻塞与max.block.ms超时异常。 - 用异步回调统计成功/失败条数,验证进程退出前不调 flush 会丢多少条。