Spring Boot 3集成Kafka:消息确认回调失效问题求助
Spring Boot 3 + Spring Kafka:消费者ACK后触发回调的实现方案
当前kafkaTemplate.send()的回调仅在消息成功发送至Kafka broker时触发,无法感知消费者处理完成并执行acknowledge()的时机。以下两种方案可实现消费者ACK后的回调通知:
方案一:使用ReplyingKafkaTemplate实现请求响应模式
该方案通过Spring Kafka原生的请求响应机制,在消费者处理完消息并ACK后,主动发送回复消息,生产者收到回复后触发回调。
1. 配置ReplyingKafkaTemplate相关Bean
@Configuration public class KafkaConfig { @Value("${spring.kafka.producer.bootstrap-servers}") private String bootstrapServers; @Bean public ProducerFactory<Object, String> producerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); return new DefaultKafkaProducerFactory<>(configProps); } @Bean public ConcurrentKafkaListenerContainerFactory<Object, String> kafkaListenerContainerFactory(ConsumerFactory<Object, String> consumerFactory) { ConcurrentKafkaListenerContainerFactory<Object, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); return factory; } @Bean public KafkaTemplate<Object, String> kafkaTemplate() { return new KafkaTemplate<>(producerFactory()); } @Bean public ReplyingKafkaTemplate<Object, String, String> replyingKafkaTemplate(ProducerFactory<Object, String> producerFactory, ConcurrentMessageListenerContainer<Object, String> replyContainer) { return new ReplyingKafkaTemplate<>(producerFactory, replyContainer); } @Bean public ConcurrentMessageListenerContainer<Object, String> replyContainer(ConsumerFactory<Object, String> consumerFactory) { ContainerProperties containerProperties = new ContainerProperties("reply-topic"); return new ConcurrentMessageListenerContainer<>(consumerFactory, containerProperties); } }
2. 修改消息发送与消费逻辑
@RestController @RequestMapping("/messages") @RequiredArgsConstructor @Slf4j public class MessageController { private final ReplyingKafkaTemplate<Object, String, String> replyingKafkaTemplate; @GetMapping public void sendMessage() { ProducerRecord<Object, String> record = new ProducerRecord<>("message-topic", "message"); // 发送消息并等待回复 RequestReplyFuture<Object, String, String> future = replyingKafkaTemplate.sendAndReceive(record); future.whenComplete((result, ex) -> { if (ex == null) { log.info("消息已被消费者处理完成,回复内容:{}", result.value()); } else { log.error("消息处理失败", ex); } }); } @KafkaListener(id = "message-topic-listener", topics = "message-topic") @SendTo("reply-topic") // 处理完成后发送回复到指定topic public String messageListener(@Payload String message, Acknowledgment acknowledgment) throws InterruptedException { Thread.sleep(10000L); log.info("处理消息 {}", message); acknowledgment.acknowledge(); return "消息处理完成:" + message; } }
3. 更新application.yml配置
添加回复消费者的组ID:
spring: kafka: producer: bootstrap-servers: localhost:29092 listener: ack-mode: MANUAL consumer: enable-auto-commit: false auto-offset-reset: earliest bootstrap-servers: localhost:29092 group-id: reply-group
方案二:自定义回调Topic实现异步通知
通过生成唯一消息标识,消费者处理完成并ACK后,向回调Topic发送标识消息,生产者监听该Topic触发对应回调逻辑,适合非阻塞场景。
1. 修改消息发送与回调监听逻辑
@RestController @RequestMapping("/messages") @RequiredArgsConstructor @Slf4j public class MessageController { private final KafkaTemplate<Object, String> kafkaTemplate; // 缓存回调逻辑,键为消息唯一标识 private final ConcurrentHashMap<String, Runnable> callbackCache = new ConcurrentHashMap<>(); @GetMapping public void sendMessage() { String correlationId = UUID.randomUUID().toString(); // 发送消息时携带唯一标识 ProducerRecord<Object, String> record = new ProducerRecord<>("message-topic", correlationId, "message"); kafkaTemplate.send(record); // 缓存回调逻辑 callbackCache.put(correlationId, () -> log.info("消息标识 {} 已被消费者处理完成", correlationId)); } @KafkaListener(id = "message-topic-listener", topics = "message-topic") public void messageListener(@Payload String message, @Header(KafkaHeaders.CORRELATION_ID) String correlationId, Acknowledgment acknowledgment) throws InterruptedException { Thread.sleep(10000L); log.info("处理消息 {}", message); acknowledgment.acknowledge(); // 发送回调通知到指定Topic kafkaTemplate.send("callback-topic", correlationId); } // 监听回调Topic触发回调 @KafkaListener(id = "callback-topic-listener", topics = "callback-topic") public void callbackListener(@Payload String correlationId) { Runnable callback = callbackCache.remove(correlationId); if (callback != null) { callback.run(); } } }
2. 更新application.yml配置
添加回调消费者的组ID:
spring: kafka: producer: bootstrap-servers: localhost:29092 listener: ack-mode: MANUAL consumer: enable-auto-commit: false auto-offset-reset: earliest bootstrap-servers: localhost:29092 group-id: main-group # 为回调监听单独配置组ID properties: spring.kafka.listener.callback-group: callback-group
方案选择
- 请求响应模式:适合需要同步等待消费者处理结果的场景,发送线程会阻塞至收到回复。
- 自定义回调Topic:适合异步通知场景,发送线程无需等待,不影响接口性能。
内容的提问来源于stack exchange,提问作者Rafael
相关产品推荐
相关产品推荐

