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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 13:45:28