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

如何在Java Kafka Producer中实现多线程发送百万级数据?

多线程Kafka生产者实现方案

嘿,我来帮你搞定这个多线程Kafka生产者的问题!你的代码已经有了基础框架,但核心问题是创建了ExecutorService却没有真正用它来分配发送任务——目前还是主线程在单线程处理所有消息发送,完全没用到多线程的优势。

先给你明确几个关键前提:

  • KafkaProducer本身是线程安全的,多个线程可以共享同一个实例来发送消息,这是官方明确支持的
  • 用ExecutorService来管理线程池,我们只需要把发送消息的逻辑包装成任务,提交给线程池即可
  • 处理百万级消息时,合理的任务提交策略可以减少线程切换开销,提升整体吞吐量

接下来是修改后的完整代码,我会标注关键修改点:

import org.apache.kafka.clients.producer.Callback;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.Metric;
import org.apache.kafka.common.MetricName;
import org.apache.kafka.common.serialization.StringSerializer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.io.BufferedReader;
import java.io.FileReader;
import java.io.IOException;
import java.util.Map;
import java.util.Properties;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;

public class KafkaProducerWithThread {
    // 初始化参数
    final String bootstrapServer = "127.0.0.1:9092";
    final String topicName = "spark-data-topic";
    final String csvFileName = "unique_products.csv";
    final static int MAX_THREAD = 4; // 建议根据CPU核心数调整,一般设为2-4倍核心数
    // 日志
    final Logger logger = LoggerFactory.getLogger(KafkaProducerWithThread.class);
    // 共享的KafkaProducer(线程安全,可多线程复用)
    private org.apache.kafka.clients.producer.KafkaProducer<String, String> producer;

    public KafkaProducerWithThread() {
        // 提前初始化生产者,避免重复创建
        this.producer = createKafkaProducer();
    }

    public static void main(String[] args) throws IOException {
        new KafkaProducerWithThread().runProducer();
    }

    public void runProducer() throws IOException {
        // 读取CSV文件
        ClassLoader classLoader = getClass().getClassLoader();
        BufferedReader reader = new BufferedReader(new FileReader(classLoader.getResource(csvFileName).getFile()));

        // 创建固定大小的线程池
        ExecutorService executorService = Executors.newFixedThreadPool(MAX_THREAD);

        String line = "";
        int i = 0;
        try {
            while ((line = reader.readLine()) != null) {
                i++;
                String key = "products_" + i;
                String value = line.trim();

                // 把发送任务包装成Runnable,提交给线程池异步执行
                executorService.submit(new KafkaSendTask(key, value));

                // 每1000条记录打印一次状态和指标
                if (i % 1000 == 0) {
                    Metric requestTotalMetric = null;
                    for (Map.Entry<MetricName, ? extends Metric> entry : producer.metrics().entrySet()) {
                        if ("request-total".equals(entry.getKey().name())) {
                            requestTotalMetric = entry.getValue();
                            break;
                        }
                    }
                    logger.info("Record #: " + i + " Request total: " + (requestTotalMetric != null ? requestTotalMetric.metricValue() : "N/A"));
                }
            }
        } finally {
            // 关闭线程池:停止接受新任务,等待已有任务完成
            executorService.shutdown();
            try {
                // 等待所有发送任务完成,最多等待5分钟(可根据实际场景调整)
                if (!executorService.awaitTermination(5, TimeUnit.MINUTES)) {
                    // 超时则强制关闭剩余任务
                    executorService.shutdownNow();
                    logger.warn("线程池超时未关闭,已强制终止");
                }
                // 确保所有未发送的消息都被flush到Kafka,再关闭生产者
                producer.flush();
                producer.close();
                logger.info("生产者已安全关闭");
            } catch (InterruptedException e) {
                executorService.shutdownNow();
                logger.error("线程池关闭被中断", e);
                // 恢复线程中断状态
                Thread.currentThread().interrupt();
            }
            reader.close();
        }
    }

    // 封装发送消息的任务类,每个任务对应一条CSV记录的发送
    private class KafkaSendTask implements Runnable {
        private final String key;
        private final String value;

        public KafkaSendTask(String key, String value) {
            this.key = key;
            this.value = value;
        }

        @Override
        public void run() {
            ProducerRecord<String, String> record = new ProducerRecord<>(topicName, key, value);
            producer.send(record, new Callback() {
                @Override
                public void onCompletion(RecordMetadata recordMetadata, Exception e) {
                    if (e != null) {
                        logger.error("发送消息失败,key: " + key, e);
                    } else {
                        // 生产环境建议关闭此日志,或改为debug级别,避免性能损耗
                        // logger.debug("消息发送成功,topic: {}, partition: {}, offset: {}",
                        //         recordMetadata.topic(), recordMetadata.partition(), recordMetadata.offset());
                    }
                }
            });
        }
    }

    public org.apache.kafka.clients.producer.KafkaProducer<String, String> createKafkaProducer() {
        Properties properties = new Properties();
        properties.setProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServer);
        properties.setProperty(ProducerConfig.ACKS_CONFIG, "all");
        properties.setProperty(ProducerConfig.RETRIES_CONFIG, "5");
        properties.setProperty(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        properties.setProperty(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        properties.setProperty(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");

        // 高吞吐量配置(以少量延迟为代价提升发送效率)
        properties.setProperty(ProducerConfig.COMPRESSION_TYPE_CONFIG, "snappy");
        properties.setProperty(ProducerConfig.LINGER_MS_CONFIG, "60");
        properties.setProperty(ProducerConfig.BATCH_SIZE_CONFIG, Integer.toString(32 * 1024)); // 32KB批量大小

        return new org.apache.kafka.clients.producer.KafkaProducer<>(properties);
    }
}

关键修改点说明:

  1. 封装发送任务:创建KafkaSendTask内部类实现Runnable,把单条消息的发送逻辑封装起来,方便提交给线程池执行
  2. 正确使用线程池:将原来主线程的发送逻辑改为executorService.submit(...),让线程池分配线程异步处理发送任务,真正利用多线程提升效率
  3. 完善资源关闭逻辑:
    • 使用executorService.shutdown()停止接受新任务,再用awaitTermination()等待所有任务完成,避免消息丢失
    • 只有所有发送任务完成后,才flush并关闭KafkaProducer,保证消息可靠性
    • 处理中断异常,确保资源能正确释放
  4. 优化指标获取:每次打印状态时重新从producer获取metrics,避免缓存过期的指标数据
  5. 调整线程池大小:根据CPU核心数设置合理的线程数,避免过多线程导致上下文切换开销

额外优化建议:

  • 批量提交任务:可以一次性读取多行CSV(比如100行),把批量消息的发送逻辑包装成一个任务提交,进一步减少线程切换次数
  • 生产者实例优化:如果线程数极大,也可以考虑给每个线程分配独立的KafkaProducer实例(会增加资源消耗),但共享实例已经能满足大部分场景的性能需求
  • 错误处理增强:可以给发送失败的消息增加重试逻辑,或者将失败消息写入死信队列,方便后续排查和重试

内容的提问来源于stack exchange,提问作者Jestino Sam

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 07:54:28