Kafka编程入门:能否用类Kafka API实现发消息后等待响应?
Kafka发送消息后等待响应的实现方案
完全可以通过Kafka原生API实现发送消息后等待响应的操作,核心是利用KafkaProducer的发送结果机制,下面是两种常用实现方式:
1. 同步发送(直接阻塞等待结果)
使用KafkaProducer.send()方法返回的Future<RecordMetadata>对象,调用其get()方法会阻塞当前线程,直到消息发送完成或抛出异常,直接拿到发送结果。
示例代码(Java):
import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; import java.util.Properties; public class SyncKafkaProducer { public static void main(String[] args) { Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) { ProducerRecord<String, String> record = new ProducerRecord<>("test-topic", "key", "hello kafka"); // 阻塞等待发送结果 RecordMetadata metadata = producer.send(record).get(); System.out.printf("消息发送成功,分区:%d,偏移量:%d%n", metadata.partition(), metadata.offset()); } catch (Exception e) { System.err.println("消息发送失败:" + e.getMessage()); e.printStackTrace(); } } }
2. 异步发送+同步等待(灵活控制等待逻辑)
如果需要更灵活的等待策略(比如设置超时时间),可以结合CountDownLatch这类同步工具,在发送回调中触发信号,主线程等待信号完成。
示例代码(Java):
import org.apache.kafka.clients.producer.*; import java.util.Properties; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; public class AsyncWaitKafkaProducer { public static void main(String[] args) { Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); CountDownLatch latch = new CountDownLatch(1); try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) { ProducerRecord<String, String> record = new ProducerRecord<>("test-topic", "key", "hello async kafka"); producer.send(record, new Callback() { @Override public void onCompletion(RecordMetadata metadata, Exception exception) { if (exception == null) { System.out.printf("消息发送成功,分区:%d,偏移量:%d%n", metadata.partition(), metadata.offset()); } else { System.err.println("消息发送失败:" + exception.getMessage()); } // 触发信号,通知主线程结束等待 latch.countDown(); } }); // 等待最多10秒,超时则结束等待 if (!latch.await(10, TimeUnit.SECONDS)) { System.err.println("消息发送等待超时"); } } catch (Exception e) { e.printStackTrace(); } } }
注意事项
- 超时控制:无论是
get()还是await(),都建议设置超时时间,避免线程无限阻塞。 - 异常处理:必须捕获发送过程中可能出现的异常(比如网络故障、主题不存在等),避免程序崩溃。
- 性能影响:同步发送会降低生产者的吞吐量,适合对消息发送结果强依赖的场景;如果追求高吞吐量,优先考虑异步发送+非阻塞的结果处理。
内容的提问来源于stack exchange,提问作者Khorrami RezA
相关产品推荐
相关产品推荐

