Producer 性能优化
定义与作用
Producer 端性能优化聚焦于三个杠杆:批量发送(减少网络往返)、数据压缩(减少传输字节)、异步发送(消除等待延迟)。正确调优可以让单 Producer 达到 100K-1M 条/秒的吞吐。
核心原理
Producer 内部流水线
两个关键缓冲:
RecordAccumulator(buffer.memory,默认 32MB):暂存待发送的消息- 每分区
Deque<ProducerBatch>:等待攒满或超时
性能三要素
配置速查表
| 参数 | 默认值 | 吞吐优先 | 延迟优先 | 说明 |
|---|---|---|---|---|
batch.size | 16384 (16KB) | 65536-131072 | 0-512 | 攒批大小 |
linger.ms | 0 | 10-100 | 0 | 等多久再发 |
buffer.memory | 33554432 (32MB) | 67108864+ | — | 总缓冲区 |
compression.type | none | lz4 / zstd | none | 压缩算法 |
max.in.flight.requests.per.connection | 5 | 5 | 1 | 并发请求数 |
acks | all (3.x) | 1 | 0 | 确认级别 |
enable.idempotence | true (3.x) | true | false | 幂等 |
完整示例
示例一:高吞吐配置(日志收集)
Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092");
props.put("key.serializer", "...");
props.put("value.serializer", "...");
// 高吞吐核心
props.put("batch.size", 131072); // 128KB 攒批
props.put("linger.ms", 50); // 等 50ms
props.put("compression.type", "lz4"); // LZ4 压缩
props.put("buffer.memory", 134217728); // 128MB
props.put("max.in.flight.requests.per.connection", 5);
测试结果:
bin/kafka-producer-perf-test.sh --topic perf-test --num-records 5000000 \
--record-size 100 --throughput -1 --producer-props \
bootstrap.servers=localhost:9092 batch.size=131072 linger.ms=50 compression.type=lz4
# 结果
# 5000000 records sent, 250000 records/sec, 23.84 MB/sec
# avg latency: 15 ms, max latency: 120 ms
示例二:低延迟配置(实时风控)
Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092");
// 低延迟核心
props.put("batch.size", 0); // 不攒批
props.put("linger.ms", 0); // 立即发送
props.put("acks", "1"); // 快速确认
props.put("max.in.flight.requests.per.connection", 1); // 保序
测试结果:
# 低延迟配置
bin/kafka-producer-perf-test.sh --topic perf-test --num-records 100000 \
--record-size 100 --throughput -1 --producer-props \
bootstrap.servers=localhost:9092 batch.size=0 linger.ms=0 acks=1
# 结果
# 100000 records sent, 12000 records/sec, 1.14 MB/sec
# avg latency: 2 ms, max latency: 8 ms
操作前后对比:
| 指标 | 默认配置 | 高吞吐配置 | 低延迟配置 |
|---|---|---|---|
| 每秒消息数 | ~50K | ~250K | ~12K |
| 平均延迟 | ~30ms | ~15ms | ~2ms |
| P99 延迟 | ~200ms | ~120ms | ~8ms |
| 网络效率 | 中 | 高(压缩+攒批) | 低(频繁小包) |
易错场景
易错 1:linger.ms=0 + 极高写入速率导致"小包轰炸"
场景:linger.ms=0 但每秒写入 100 万条。
后果:每条消息一个独立的 ProduceRequest → 大量小数据包 → 网络带宽利用率 < 30%,CPU 大量用于网络中断处理。
解决:吞吐场景永远设置 linger.ms=5 以上。
易错 2:buffer.memory 太小导致 send() 阻塞
场景:buffer.memory=32MB(默认),写入速率 > Broker 处理速率。
后果:RecordAccumulator 被写满 → send() 阻塞最长达 max.block.ms(默认 60s)→ 应用线程卡住。
诊断:
// 监控缓冲区使用率
float usage = (float) producer.metrics()
.get("buffer-total-bytes-used").metricValue()
/ producer.metrics().get("buffer-total-bytes-max").metricValue();
面试高频考点
Q:batch.size 和 linger.ms 的关系是什么?谁先触发就发送?
A:
batch.size是空间触发:攒够指定字节就发linger.ms是时间触发:超过指定毫秒就发- 先满足谁就触发:如果 5ms 攒够了 128KB → 以
batch.size触发;如果 10ms 还没攒够 128KB → 以linger.ms触发
实际调优中,linger.ms 是更主要的控制变量——它保证了延迟的确定性上限。