Kafka多线程生产者消息丢失问题求助(Kafka 1.1.0)
Kafka消息丢失问题分析与修复方案
嘿,咱们来看看为啥你的Kafka集群里丢了大概5000条消息。核心问题在于你没正确处理Kafka生产者的异步特性,还漏掉了关键的资源清理步骤。咱们一步步拆解并解决这个问题:
核心问题点
1. 异步发送未等待确认,计数不准
producer.send()是异步非阻塞方法——调用后立刻返回,但此时消息可能还在生产者本地缓冲区躺着,或者正在网络传输路上。你的代码在调用send()后立刻递增count,但这个数值只代表你发起了多少次发送请求,完全不代表消息已经被Kafka集群成功接收。当线程循环结束后,生产者实例被销毁,缓冲区里没发出去的消息直接就丢了。
2. 没调用flush()和close()确保消息全部提交
每个线程创建的KafkaMessageSender在循环结束后,既没调用producer.flush()(强制把缓冲区里的消息推送到集群),也没调用producer.close()(优雅关闭生产者,等所有未完成的发送请求处理完)。这直接导致大量滞留在缓冲区的消息被丢弃,这就是你丢了几千条的主要原因。
3. 不必要的多生产者实例(次要问题)
KafkaProducer是线程安全的,完全可以让所有线程共享同一个实例,不用每个线程都新建一个。多实例不仅浪费资源,还会因为每个实例的缓冲区配置,导致更多消息滞留。
修复后的代码示例
修改KafkaMessageSender类
public class KafkaMessageSender { private final static Logger logger = LoggerFactory.getLogger(KafkaMessageSender.class); private final KafkaProducer<String, String> producer; private final String topic; private final AtomicInteger successCount; private final AtomicInteger failCount; public KafkaMessageSender(KafkaProducer<String, String> producer, String topic, AtomicInteger successCount, AtomicInteger failCount) { logger.info("KafkaMessageSender initializing..."); this.producer = producer; this.topic = topic; this.successCount = successCount; this.failCount = failCount; logger.info("KafkaMessageSender initializing end"); } public void sendMessages() { ProducerRecord<String, String> record = new ProducerRecord<>(topic, Messages.MSG_4K); // 添加回调,真正统计成功/失败的消息数 producer.send(record, (metadata, exception) -> { if (exception != null) { logger.error("Message send failed", exception); failCount.getAndIncrement(); } else { successCount.getAndIncrement(); // 可选:打印成功的offset,方便排查 // logger.info("Message sent successfully, offset: {}", metadata.offset()); } }); } // 添加flush方法,强制推送缓冲区消息 public void flush() { producer.flush(); } }
修改KafkaMessageSenderMain类
public class KafkaMessageSenderMain { private final static Logger logger = LoggerFactory.getLogger(KafkaMessageSenderMain.class); final static String bootstrap_url = "ism1.solulink.co.kr:9092,ism2.solulink.co.kr:9092,ism3.solulink.co.kr:9092"; final static String topic = "test-topic"; final static AtomicInteger successCount = new AtomicInteger(0); final static AtomicInteger failCount = new AtomicInteger(0); final static int MAX_LOOP = 10000; final static int MAX_THREAD = 5; public static void main(String[] args) { long startTime = System.currentTimeMillis(); // 全局创建一个生产者实例,所有线程共享 Properties props = new Properties(); props.put(ProducerConfig.ACKS_CONFIG, "all"); // 注意:你原代码重复设置了ACKS,这里保留最终生效的"all" props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrap_url); props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); props.put(ProducerConfig.RETRIES_CONFIG, 3); props.put(ProducerConfig.LINGER_MS_CONFIG, 5); props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); KafkaProducer<String, String> sharedProducer = new KafkaProducer<>(props); ExecutorService executorService = Executors.newFixedThreadPool(MAX_THREAD); for(int i = 0; i < MAX_THREAD; i++) { executorService.execute(() -> { KafkaMessageSender sender = new KafkaMessageSender(sharedProducer, topic, successCount, failCount); for(int j = 0; j < MAX_LOOP; j++) { sender.sendMessages(); } // 线程循环结束后,先flush确保缓冲区消息都被推送 sender.flush(); }); } executorService.shutdown(); try { boolean flag = executorService.awaitTermination(Long.MAX_VALUE, TimeUnit.SECONDS); // 所有线程结束后,优雅关闭生产者,等待所有未完成的请求 sharedProducer.close(); long endTime = System.currentTimeMillis(); long procTime = (endTime - startTime); logger.info("All threads are shutdown? : {}", flag); logger.info("Process time : {} sec", (double)procTime/1000L); logger.info("Total success messages : {}", successCount.get()); logger.info("Total failed messages : {}", failCount.get()); } catch (InterruptedException e) { logger.error("Await termination exception", e); // 中断时也要关闭生产者,避免消息丢失 sharedProducer.close(); } } }
额外提醒
- ACK配置:你原代码里重复设置了
ACKS_CONFIG(先设"1"又设"all"),最终生效的是"all",这是正确的选择——"all"要求所有副本都确认接收,能最大程度保证消息不丢失。 - 回调的作用:添加回调后,你能精准统计真正成功/失败的消息数,而不是依赖发起请求的次数,这能帮你更准确地排查问题。
- 生产者生命周期:
producer.close()会阻塞直到所有未完成的发送请求都得到集群确认,一定要在所有发送任务完成后再调用它。
按照这个修改后,你应该能看到实际的offset接近50000,消息丢失的问题就能解决啦。
内容的提问来源于stack exchange,提问作者dev_jongtae
相关产品推荐
相关产品推荐

