要让kafka生产者既不丢消息也不重复发送,核心是重试机制与幂等性协同:需设delivery.timeout.ms(如120000)兜底重试时长,enable.idempotence=true且满足acks=all、max.in.flight≤5、retries>0三条件,并辅以linger.ms与batch.size优化批量发送。

要让 Kafka 生产者既不丢消息、也不重复发送,核心是把重试机制和幂等性配合好。单独开重试可能重复,只开幂等不配重试又可能失败丢数据——两者必须协同设置。
重试配置:别只设 retries,重点看 delivery.timeout.ms
retries 控制最大重试次数,但真正决定“重试多久”的是 delivery.timeout.ms(默认 2 分钟)。Kafka 2.1+ 默认 retries=Integer.MAX_VALUE,所以实际重试行为由这个超时值兜底。
- 建议显式设置 delivery.timeout.ms=120000(2 分钟),避免无限重试拖垮线程或积压内存
- 不要设 retries=0 —— 这等于放弃重试,ack=all 也救不回网络抖动导致的瞬时失败
- 若业务对延迟极敏感(如实时风控),可调低到 30000(30 秒),但需接受少量不可恢复失败
幂等性开启:三要素缺一不可
enable.idempotence=true 不是“开了就完事”,它依赖三个硬性条件才能生效:
- acks 必须为 all(或 -1)—— 否则服务端无法保证 ISR 副本同步完成,sequence 校验会失效
- max.in.flight.requests.per.connection ≤ 5(Kafka ≥ 1.1)—— 超过会引发 OutOfOrderSequenceException;推荐设为 5,兼顾吞吐与顺序
- retries > 0(哪怕只设为 1)—— 幂等性只对重试场景起作用,零重试下它不生效
配置示例:
props.put("enable.idempotence", "true");props.put("acks", "all");
props.put("max.in.flight.requests.per.connection", "5");
props.put("delivery.timeout.ms", "120000");
补充加固:linger.ms + batch.size 提升稳定性
批量发送能减少网络请求压力,间接降低重试概率。但 linger.ms 设太高会增加端到端延迟,太低则失去批量意义。
- 常规业务推荐 linger.ms=5–20,batch.size=16384(16KB)
- 高吞吐日志类场景可设 linger.ms=100,但需监控 RecordAccumulator 内存占用
- 注意:batch 太大 + linger 太长,可能在超时前积压过多消息,触发 delivery.timeout 中断重试
为什么还要消费者端做幂等?
生产者幂等只保障“单 Producer、单分区、单会话内不重复”。一旦 Producer 实例重启、扩容缩容、或跨集群迁移,PID 重置,sequence 归零——历史重复风险依然存在。
- 业务关键字段(如订单号、支付流水号)加唯一索引,数据库层拦截重复插入
- 消费逻辑中用 Redis 记录已处理 msgId(带 TTL),先查后执行
- 避免仅依赖 offset 提交来防重——自动提交或 rebalance 时仍可能漏处理
Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南











