如何在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); } }
关键修改点说明:
- 封装发送任务:创建
KafkaSendTask内部类实现Runnable,把单条消息的发送逻辑封装起来,方便提交给线程池执行 - 正确使用线程池:将原来主线程的发送逻辑改为
executorService.submit(...),让线程池分配线程异步处理发送任务,真正利用多线程提升效率 - 完善资源关闭逻辑:
- 使用
executorService.shutdown()停止接受新任务,再用awaitTermination()等待所有任务完成,避免消息丢失 - 只有所有发送任务完成后,才flush并关闭KafkaProducer,保证消息可靠性
- 处理中断异常,确保资源能正确释放
- 使用
- 优化指标获取:每次打印状态时重新从producer获取metrics,避免缓存过期的指标数据
- 调整线程池大小:根据CPU核心数设置合理的线程数,避免过多线程导致上下文切换开销
额外优化建议:
- 批量提交任务:可以一次性读取多行CSV(比如100行),把批量消息的发送逻辑包装成一个任务提交,进一步减少线程切换次数
- 生产者实例优化:如果线程数极大,也可以考虑给每个线程分配独立的KafkaProducer实例(会增加资源消耗),但共享实例已经能满足大部分场景的性能需求
- 错误处理增强:可以给发送失败的消息增加重试逻辑,或者将失败消息写入死信队列,方便后续排查和重试
内容的提问来源于stack exchange,提问作者Jestino Sam
相关产品推荐
相关产品推荐

