Producer 发送机制
定义与作用
Producer 发送机制描述了消息从 send() 调用到 Broker 确认写入的完整链路。这条链路决定了消息的可靠性和延迟——理解它,才能正确配置 ACK、重试和批处理参数。
核心原理
完整发送流程
ACK 机制的三级保证
acks 是 Producer 最重要的可靠性参数:
| acks 值 | 发送延迟 | 可靠性 | 适用场景 |
|---|---|---|---|
0 | 最低 | 消息可能丢失 | 指标监控(丢几条无所谓) |
1 | 中等 | Leader 崩溃时可能丢失 | 日志收集(可容忍少量丢失) |
all / -1 | 最高 | 不丢失(配合 min.insync.replicas) | 交易、订单等关键数据 |
重试机制
重要:retries 配合 enable.idempotence=true 使用时,Kafka 会自动将 max.in.flight.requests.per.connection 限制为 5,保证消息顺序。
完整示例
示例一:三种 ACK 策略的吞吐对比
场景:测试 acks 对吞吐量的影响。
操作前环境:Topic ack-test,3 分区,3 副本。
# acks=0
bin/kafka-producer-perf-test.sh --topic ack-test --num-records 100000 \
--record-size 100 --throughput -1 \
--producer-props acks=0 linger.ms=5 bootstrap.servers=localhost:9092
# 100000 records sent, 142857 records/sec (13.62 MB/sec)
# acks=1
bin/kafka-producer-perf-test.sh --topic ack-test --num-records 100000 \
--record-size 100 --throughput -1 \
--producer-props acks=1 linger.ms=5
# 100000 records sent, 89285 records/sec (8.51 MB/sec)
# acks=all
bin/kafka-producer-perf-test.sh --topic ack-test --num-records 100000 \
--record-size 100 --throughput -1 \
--producer-props acks=all linger.ms=5
# 100000 records sent, 35714 records/sec (3.41 MB/sec)
操作后对比:
| acks | 吞吐量 | 延迟(平均) | 可靠性 |
|---|---|---|---|
| 0 | 142,857/s | 0.05ms | 可能丢失 |
| 1 | 89,285/s | 1.2ms | Leader 故障可能丢失 |
| all | 35,714/s | 3.8ms | 不丢(含副本) |
示例二:配置重试避免消息丢失
场景:网络抖动导致 Producer 发送失败,需要自动重试。
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// ==== 可靠性优先配置 ====
props.put("acks", "all"); // 等待所有 ISR
props.put("retries", Integer.MAX_VALUE); // 无限重试(实际上由 delivery.timeout.ms 控制)
props.put("max.in.flight.requests.per.connection", 5); // 幂等时默认5,防止乱序
props.put("enable.idempotence", true); // 开启幂等,防止重复
props.put("delivery.timeout.ms", 120000); // 2分钟超时
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
ProducerRecord<String, String> record = new ProducerRecord<>("orders", "order-1", "data");
producer.send(record, (metadata, exception) -> {
if (exception != null) {
// delivery.timeout.ms 耗尽后才会进入这里
System.err.println("最终发送失败: " + exception.getMessage());
}
});
}
操作前后对比:
| 配置 | 无重试(retries=0) | 有重试(retries=MAX) |
|---|---|---|
| 网络闪断 100ms | 消息丢失 | 自动重试成功 |
| Leader 切换 | 消息丢失 | 刷新元数据后重试成功 |
| 连续故障 >2min | 消息丢失 | delivery.timeout.ms 超时后失败 |
易错场景
易错 1:acks=all 但 min.insync.replicas=1
场景:Producer 配置 acks=all,但没有设置 Broker 端的 min.insync.replicas。
后果:当 ISR 只有 Leader 一个节点时(Follower 故障),acks=all 退化为 acks=1,不具备真正的可靠性。
正确做法:Topic 级别设置 min.insync.replicas=2,确保至少 2 个副本在 ISR 中才允许写入。
易错 2:max.in.flight.requests.per.connection 与顺序性
场景:设置 max.in.flight.requests.per.connection=5(允许 5 个未确认请求并发),且未开启幂等。
后果:如果 batch1 发送失败被重试,而 batch2-batch5 已成功写入,重试成功后的 batch1 会写在 batch5 之后——消息乱序。
规则:
- 需要严格有序且未开启幂等 →
max.in.flight.requests.per.connection=1 - 开启幂等 → Kafka 自动限制为 5,且保证顺序
面试高频考点
Q:acks=all 是否绝对保证消息不丢失?
A:不绝对。acks=all 保证的是"消息被 ISR 中所有副本确认后才认为发送成功"。但如果:
- ISR 中的所有副本在确认后、返回 Producer 响应前的瞬间全部同时崩溃(极小概率)
min.insync.replicas=1且 ISR 只有 Leader,退化为acks=1
真正的"绝对不丢"需要在 Producer(acks=all + retries=MAX + 幂等)、Broker(min.insync.replicas >= 2 + unclean.leader.election.enable=false)和 Consumer(手动提交 Offset、处理完再提交)三个层面同时保证。