Kafka生产者异步模式疑问:两次send调用未按预期异步执行
Kafka异步send方法为何出现阻塞执行的问题
问题核心
当Kafka服务关闭时,连续调用两次异步send方法,预期会立刻输出两条Thread before,但实际第一条send触发超时异常后,才会执行第二条send的前置输出。
代码复现
消息发送代码
private static void send(KafkaProducer<String, String> producer) { System.out.println(Thread.currentThread()+" before"); producer.send(new ProducerRecord<String, String>("test", "test"), (metadata, exception) -> { if (exception == null) { System.out.println(metadata.partition()); } else { exception.printStackTrace(); } }); System.out.println(Thread.currentThread()+" after"); }
Kafka配置代码
Properties properties = new Properties(); properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9091"); properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); KafkaProducer<String, String> producer = new KafkaProducer<>(properties);
调用方式
send(producer); send(producer);
预期与实际输出
- 预期输出:
Thread before Thread before
- 实际输出:
Thread[main,5,main] before org.apache.kafka.common.errors.TimeoutException: Topic test not present in metadata after 60000 ms. Thread[main,5,main] after Thread[main,5,main] before
原因分析
KafkaProducer的send方法并非完全无阻塞:
- 第一次调用
send时,由于Kafka服务未启动,producer需要获取目标topic的元数据,但无法连接到broker,会触发同步元数据拉取流程。 - 这个元数据拉取过程会阻塞主线程,直到达到默认的元数据获取超时时间(60秒),才会抛出
TimeoutException并继续执行后续代码。 - 因此主线程卡在第一次
send调用中,直到超时后才会执行第二个send的Thread before输出。
解决方案
如果想要实现预期的并行触发效果,可以通过以下方式:
- 多线程调用send:将每次
send的调用放到独立的线程中,避免主线程被元数据拉取阻塞:
new Thread(() -> send(producer)).start(); new Thread(() -> send(producer)).start();
- 调整元数据超时时间:通过
ProducerConfig.METADATA_FETCH_TIMEOUT_CONFIG缩短超时时间,减少阻塞时长,但无法完全消除阻塞:
properties.put(ProducerConfig.METADATA_FETCH_TIMEOUT_CONFIG, 1000); // 设置为1秒超时
内容的提问来源于stack exchange,提问作者Violetta
相关产品推荐
相关产品推荐

