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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 21:00:05