You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.28 09:45:05