Producer 配置
定义与作用
本节提供 Producer 关键配置参数的速查表与调优指南。不罗列所有参数,聚焦于影响吞吐、延迟、可靠性的核心参数及其组合策略。
核心原理
Producer 的配置参数关系着三个维度的权衡:
配置速查表
核心连接与发送
| 参数 | 默认值 | 说明 | 调优建议 |
|---|---|---|---|
bootstrap.servers | (必填) | Broker 地址列表 | 填写多个,防止单点连接失败 |
key.serializer | (必填) | Key 序列化器 | StringSerializer / ByteArraySerializer |
value.serializer | (必填) | Value 序列化器 | 同上 |
client.id | "" | 客户端标识 | 建议设置,便于日志追踪 |
批处理与吞吐
| 参数 | 默认值 | 说明 | 调优建议 |
|---|---|---|---|
batch.size | 16384 (16KB) | 每个分区的批次大小 | 吞吐优化往上调至 32KB-512KB |
linger.ms | 0 | 批次等待时间(ms) | 吞吐优化设 5-100;延迟敏感保持 0 |
buffer.memory | 33554432 (32MB) | 缓冲区总大小 | 高吞吐场景调至 64-256MB |
compression.type | none | 压缩类型(none/gzip/snappy/lz4/zstd) | 推荐 lz4(吞吐/CPU 比最优)或 zstd(高压缩率) |
max.request.size | 1048576 (1MB) | 单次请求最大字节 | 大数据场景上调,需同步调整 Broker message.max.bytes |
可靠性
| 参数 | 默认值 | 说明 | 调优建议 |
|---|---|---|---|
acks | all(3.0+) | 确认级别 | 关键数据必须 all |
enable.idempotence | true(3.0+) | 幂等生产者 | 建议保持开启 |
retries | 2147483647 | 重试次数 | 与 delivery.timeout.ms 配合 |
delivery.timeout.ms | 120000 (2min) | 发送超时(含重试) | 根据 SLA 调整 |
retry.backoff.ms | 100 | 重试间隔(ms) | 网络抖动频繁可适当增大 |
max.in.flight.requests.per.connection | 5 | 未确认请求数 | 幂等时自动 5;不幂等需有序时设为 1 |
事务
| 参数 | 默认值 | 说明 |
|---|---|---|
transactional.id | null | 事务 ID(非空即开启事务模式) |
transaction.timeout.ms | 60000 (1min) | 事务超时 |
完整示例
示例一:高吞吐日志收集配置
场景:日志收集系统,每秒 50 万条小消息(每条 ~200B),可容忍少量丢失。
Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092,kafka3:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// 吞吐优先
props.put("acks", "1");
props.put("linger.ms", 10);
props.put("batch.size", 131072); // 128KB
props.put("buffer.memory", 134217728); // 128MB
props.put("compression.type", "lz4");
props.put("max.in.flight.requests.per.connection", 5);
操作前后对比:
| 配置 | 默认配置 | 高吞吐配置 |
|---|---|---|
linger.ms | 0 | 10 |
batch.size | 16KB | 128KB |
compression.type | none | lz4 |
| 实测吞吐 | ~50K/s | ~200K/s |
示例二:交易系统可靠性优先配置
场景:交易系统,要求消息绝对不丢失、不重复、不乱序。
Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092,kafka3:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// 可靠性优先
props.put("enable.idempotence", true); // acks=all, retries=MAX, max.in.flight=5
props.put("delivery.timeout.ms", 120000); // 2 分钟超时
props.put("transactional.id", "tx-trade-001"); // 开启事务(如需跨 Topic 原子性)
注意:可靠性优先配置配合 Broker 端:
- Topic 级别:
min.insync.replicas >= 2 - Broker 级别:
unclean.leader.election.enable=false
易错场景
易错 1:compression.type 设为 gzip 导致延迟飙升
场景:延迟敏感系统(P99 < 10ms)使用 gzip 压缩。
后果:gzip CPU 消耗高,压缩/解压时间可能超过网络传输时间,延迟不降反升。
建议:
- 延迟优先:
snappy或none - 吞吐优先:
lz4(CPU/压缩比最佳平衡) - 存储优先:
zstd(最高压缩比)
易错 2:buffer.memory 太小导致 block.on.buffer.full
场景:高吞吐场景下 buffer.memory=32MB(默认)。
后果:缓冲区满后 send() 阻塞,等待空间释放(受 max.block.ms 控制,默认 60s)。
诊断:
// 监控 Producer 指标
// bufferpool-wait-time: 等待缓冲区的时间(应为 0)
// buffer-available-bytes: 缓冲区可用字节(不应经常为 0)
面试高频考点
Q:linger.ms=0 时 batch.size 还有意义吗?
A:有意义。linger.ms 和 batch.size 是"或"的关系——哪个条件先达到就发送:
linger.ms=0意味着到达即发- 但如果前一个请求还在传输中,新消息会继续积累
- 积累到
batch.size上限时,即使 linger.ms 未到也会发送
所以 linger.ms=0 + 高并发场景下,批次仍然会填满 batch.size(因为总有消息在排队等待前一个请求完成)。linger.ms 的主要作用是在低负载场景下多等一会让批次填满。