
一个被忽视的 Kafka 生产者模型细节引发了生产环境的静默丢数据事故。一、事故现场某天凌晨日志推送服务开始频繁告警org.springframework.kafka.core.KafkaProducerException: Failed to send; nested exception is org.apache.kafka.common.errors.TimeoutException: Expiring 120 record(s) for dmp-biz-logs-3: 120609 ms has passed since batch creation多个 topic、多个分区同时报错。更糟糕的是由于 Kafka 发送是异步的HTTP 接口早已返回 200调用方毫不知情——日志在链路中静默消失了。二、追根溯源2.1 当时的架构日志入口是一个 Spring Boot 服务核心代码长这样ServicepublicclassCommonLogServiceImpl{AutowiredIKafkaServicekafkaService;// → 包装了 KafkaTemplatepublicvoidprocessLogs(Stringtopic,Stringmsg){JsonNodejsonNodeobjectMapper.readTree(msg);if(jsonNode.isArray()){IteratorJsonNodeiteratorjsonNode.iterator();while(iterator.hasNext()){JsonNodenodeiterator.next();// 逐条发送到 KafkakafkaService.producer(topic,node.toString());// ← 关键行}}}}KafkaServiceImpl里是标准的KafkaTemplate.send()ServicepublicclassKafkaServiceImplimplementsIKafkaService{AutowiredprivateKafkaTemplatebyte[],byte[]kafkaTemplate;// ← Spring Boot 自动注入的单例Overridepublicvoidproducer(Stringtopic,Stringmsg){ProducerRecordbyte[],byte[]recordnewProducerRecord(topic,msg.getBytes(StandardCharsets.UTF_8));kafkaTemplate.send(record);}}看起来没什么问题对吧异步发送、非阻塞、回调记日志标准的 Spring Kafka 写法。2.2 问题到底在哪问题藏在 Kafka 客户端的一个基本事实里每个KafkaProducer实例只有一个Sender线程。不管你调用多少次send()底层架构是这样的HTTP线程1 ──send()──┐ HTTP线程2 ──send()──┤──▶ 同一个 RecordAccumulator(32MB) ──▶ Sender × 1 ──▶ Kafka Broker HTTP线程3 ──send()──┘ ↑ ↑ HTTP线程N ──send()──┘ 所有消息往一个缓冲区塞 只有一个线程真正发而KafkaTemplate是 Spring 容器管理的单例 Bean整个 JVM 进程只有一个KafkaProducer只有一个Sender。当/handleLogEvents接口达到410 req/s每个请求携带几十上百条日志所有消息涌向同一个RecordAccumulator默认 32MB。Sender 线程需要依次完成序列化、压缩、网络发送、等待 Broker 确认——它根本消化不过来。灾难链条如下① 消息涌入速度 Sender 消费速度 ↓ ② RecordAccumulator 堆积队尾消息排队越来越久 ↓ ③ 排队超过 120 秒delivery.timeout.ms 默认值 ↓ ④ Kafka 客户端判定消息过期丢弃 抛 TimeoutException ↓ ⑤ 但 HTTP 接口早已返回 200 —— 数据静默丢失更糟的是当缓冲区满了32MB 打满后续send()会阻塞等待空位max.block.ms默认 60 秒直接把 HTTP 线程拖死。三、解法对比方案一调参缓解不根治改 Kafka Producer 参数让每条消息等得更短、批次更高效spring.kafka.producer.properties.linger.ms5 # 微批聚合 5ms减少 Sender 处理次数 spring.kafka.producer.properties.batch.size65536 # 增大批次 spring.kafka.producer.properties.max.block.ms5000 # 等 5 秒进不去就放弃不拖死 HTTP优点不改代码Apollo 动态下发即时生效。缺点流量再涨还是会复发Sender 单线程的天花板没变。方案二加内存队列解耦改动大在 HTTP 层和 Kafka 层之间插入内存队列HTTP 线程只负责入队非阻塞后台线程按自己节奏消费发送HTTP线程 ──offer(100ms)──▶ 内存队列 ──▶ 后台线程 ──▶ Kafka ↑ 永不阻塞 ↑ 削峰缓冲 ↑ 按自己节奏优点彻底解耦HTTP 永不阻塞。缺点需要新建队列消费线程、处理背压、优雅关闭改动较大。方案三多 KafkaProducer 实例突破并发瓶颈改动最小既然单 Producer 单 Sender 是天花板那就用多个 Producer改造前 HTTP线程 ──send()──▶ 1 个 KafkaProducer ── 1 个 Sender ──▶ Kafka瓶颈 改造后 HTTP线程 ──send()──▶ Producer1 ── Sender1 ──▶ Kafka ──send()──▶ Producer2 ── Sender2 ──▶ Kafka ──send()──▶ Producer3 ── Sender3 ──▶ Kafka ──send()──▶ Producer4 ── Sender4 ──▶ Kafka ↑ 4 个独立的 Sender 并行发送吞吐 ×4优点改动最小新增 1 个类 改 4 行调用直接命中根因。缺点仍然是同步耦合极端流量下 send() 仍可能阻塞。四、方案三的实现4.1 KafkaProducerPoolComponentpublicclassKafkaProducerPool{privatestaticfinalLoggerloggerLoggerFactory.getLogger(KafkaProducerPool.class);AutowiredprivateKafkaPropertieskafkaProperties;Value(${kafka.producer.pool.size:4})privateintpoolSize;privatefinalListKafkaProducerbyte[],byte[]producersnewArrayList();privatefinalAtomicIntegerroundRobinnewAtomicInteger(0);PostConstructpublicvoidinit(){MapString,ObjectconfigskafkaProperties.buildProducerProperties();StringbaseClientId(String)configs.getOrDefault(client.id,kafka-producer-pool);for(inti0;ipoolSize;i){configs.put(client.id,baseClientId-i);producers.add(newKafkaProducer(configs));}logger.info(KafkaProducerPool initialized: poolSize{},poolSize);}/** * Round-Robin 轮询获取 Producer保证负载均匀 */privateKafkaProducerbyte[],byte[]getProducer(){intidxMath.abs(roundRobin.getAndIncrement()%producers.size());returnproducers.get(idx);}/** * 异步发送不阻塞调用线程 */publicvoidsend(Stringtopic,Stringmessage){ProducerRecordbyte[],byte[]recordnewProducerRecord(topic,message.getBytes(StandardCharsets.UTF_8));getProducer().send(record,(metadata,exception)-{if(exception!null){logger.error(Kafka send failed: topic{},topic,exception);}});}PreDestroypublicvoiddestroy(){producers.forEach(p-{try{p.close();}catch(Exceptione){logger.error(close error,e);}});}}核心思路就三点复用 Spring Boot 的KafkaProperties配置与原来的KafkaTemplate完全一致仅client.id加上-0/-1/-2/-3后缀区分AtomicInteger轮询分配无锁均匀异步 send 回调记日志与原行为一致4.2 调用方改动// 改前kafkaService.producer(topic,objectNode.toString());// 改后producerPool.send(topic,objectNode.toString());整个CommonLogServiceImpl只改 4 行。4.3 Apollo 配置# Producer 池大小默认 4建议 CPU 核心数 × 2 kafka.producer.pool.size4五、效果维度改造前改造后Sender 线程数1N池大小RecordAccumulator 总容量32MB32MB × N峰值吞吐受单线程限制线性增长N4 时约 ×4改动量—新增 1 个类 改 4 行外部依赖—零仅 JDKListAtomicInteger六、但是这还不够回头再看KafkaServiceImpl里面还有这个方法Overridepublicvoidproducer(Stringtopic,Stringob_object_id,ListMapString,ObjectmsgList){for(MapString,Objectmap:msgList){StringstrConstants.objectMapper.writeValueAsString(map);ProducerRecordbyte[],byte[]recordnewProducerRecord(topic,str.getBytes(StandardCharsets.UTF_8));kafkaTemplate.send(record);// ← 还是走的同一个 kafkaTemplate}}kafkaTemplate仍然是单例这个方法里的所有send()还是往同一个 Producer 同一个 Sender里塞。改完CommonLogServiceImpl只是把最大的入口解了但其他调用方依旧共用那个单 Sender隐患还在。七、终极方案Producer 池 内存队列 组合两个方案互补HTTP线程 ──offer(100ms)──▶ 内存队列 ──▶ 消费线程池 ──▶ Producer 池 ──▶ Kafka ↑ ↑ ↑ ↑ 非阻塞入队 削峰缓冲 多线程消费 多 Sender 并行方案解决的问题未解决的问题仅内存队列HTTP 解耦、削峰消费端仍是单 Sender仅多 ProducerSender 并发瓶颈HTTP 与 Kafka 仍耦合两者组合全部无明显短板这也是我们最终上线的方案。八、总结这次问题的根因不是什么高深的分布式理论而是Kafka 客户端一个容易忽略的模型细节KafkaTemplate是单例 → 只有一个KafkaProducer→ 只有一个Sender线程。在高并发场景下这个单线程就是整个系统的阿喀琉斯之踵。解决思路也很直接一个不够就用多个。但别止步于此——多 Producer 解决了并发瓶颈但没解决同步耦合。真正的生产级方案应该是解耦 并行内存队列把 HTTP 和 Kafka 隔开多 Producer让发送端不再有单点瓶颈两条腿走路才走得稳。