如何让Spring EmbeddedKafka生产者等待消费者手动确认?
解决方案:让生产者阻塞直到消费者手动确认
你的核心问题在于**acks=all仅保证消息被Kafka Broker的所有ISR副本持久化完成**,和消费者的手动确认(Acknowledgement.acknowledge())完全是两个独立环节,所以调用producer.send().get()只会等待Broker的确认,不会等待消费者处理并完成ack。
要实现测试中生产者阻塞直到消费者完成确认且无需sleep,推荐用同步计数器(CountDownLatch) 来实现,具体步骤如下:
1. 在测试类中定义CountDownLatch
在测试类里初始化一个计数器,计数设为1(匹配你配置的max-poll-records: 1,每次处理一条消息):
private CountDownLatch consumerAckLatch = new CountDownLatch(1);
2. 修改消费者方法,完成ack后触发计数器
在@KafkaListener方法中,调用Acknowledgement.acknowledge()之后,触发计数器的countDown()方法:
@KafkaListener(topics = "your-topic") public void listen(ConsumerRecord<String, MessageType> record, Acknowledgement ack) { // 执行消息处理逻辑 // ... // 手动确认消息 ack.acknowledge(); // 触发计数器,通知生产者已完成确认 consumerAckLatch.countDown(); }
3. 生产者发送后等待计数器完成
在测试代码中,发送消息后调用latch.await(),阻塞直到消费者完成ack触发计数器:
// 发送消息并等待Broker确认 producer.send(new ProducerRecord<Integer, String>(TOPIC, key, value)).get(); // 阻塞等待消费者完成手动确认,可设置超时避免无限阻塞 consumerAckLatch.await(10, TimeUnit.SECONDS);
额外注意事项
- 如果消费者是多线程模式(
setConcurrency(numListeners)值大于1),需要根据并发数调整计数器初始值,比如并发数为3时,初始值设为3。 - 每次测试前要重置计数器,避免不同测试用例互相干扰:
@BeforeEach void setUp() { consumerAckLatch = new CountDownLatch(1); }
内容的提问来源于stack exchange,提问作者user2654096
相关产品推荐
相关产品推荐

