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

无需@KafkaListener与@SendTo实现ReplyingKafkaTemplate响应发送

ReplyingKafkaTemplate 自定义监听器请求响应超时问题解决方案

核心问题分析

你当前的自定义监听器缺少请求响应的关键逻辑:ReplyingKafkaTemplate 发送请求时会自动在消息头中添加 reply-topic 和 correlation-id,消费端必须读取这两个头信息,将响应消息发送到指定的回复主题,并携带相同的 correlation-id,否则 ReplyingKafkaTemplate 无法匹配到对应的请求,最终导致超时。

具体解决方案步骤

1. 修改自定义监听器,添加响应发送逻辑

在 onMessage 方法中读取请求消息头的回复主题和关联ID,构造响应并发送:

public class KafkaListener implements AcknowledgingMessageListener<String, String> {

    private final KafkaTemplate<String, String> kafkaTemplate;

    // 通过构造函数注入KafkaTemplate用于发送响应
    public KafkaListener(KafkaTemplate<String, String> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    @Override
    public void onMessage(ConsumerRecord<String, String> consumeRecord, Acknowledgment acknowledgment) {
        try {
            // 从请求消息头中提取回复主题和关联ID
            Header replyTopicHeader = consumeRecord.headers().lastHeader("reply-topic");
            Header correlationIdHeader = consumeRecord.headers().lastHeader("correlation-id");
            if (replyTopicHeader == null || correlationIdHeader == null) {
                // 非请求响应消息,直接提交偏移量
                acknowledgment.acknowledge();
                return;
            }

            String replyTopic = new String(replyTopicHeader.value());
            String correlationId = new String(correlationIdHeader.value());

            // 处理业务逻辑,生成响应内容
            String requestContent = consumeRecord.value();
            String responseContent = "处理完成:" + requestContent;

            // 构造响应消息,携带关联ID头
            ProducerRecord<String, String> responseRecord = new ProducerRecord<>(replyTopic, responseContent);
            responseRecord.headers().add("correlation-id", correlationId.getBytes());

            // 发送响应
            kafkaTemplate.send(responseRecord);

            // 手动提交偏移量
            acknowledgment.acknowledge();
        } catch (Exception e) {
            // 异常处理,如重试或转发死信队列
            e.printStackTrace();
        }
    }
}

2. 修正回复容器配置

回复容器必须处于运行状态才能接收响应,需修改setAutoStartup配置:

@Bean
public ConcurrentMessageListenerContainer<String, String> container(
        ConcurrentKafkaListenerContainerFactory<String, String> containerFactory) {

    ConcurrentMessageListenerContainer<String, String> repliesContainer =
            containerFactory.createContainer("kReplies");
    repliesContainer.getContainerProperties().setGroupId("repliesGroup");
    // 启动回复容器,否则无法接收响应消息
    repliesContainer.setAutoStartup(true);

    return repliesContainer;
}

3. 优化ReplyingKafkaTemplate配置

可显式设置回复超时时间,确保和生产者端的等待时间匹配:

@Bean
public ReplyingKafkaTemplate<String, String, String> replyingTemplate(
        ProducerFactory<String, String> pf,
        ConcurrentMessageListenerContainer<String, String> repliesContainer) {
    ReplyingKafkaTemplate<String, String, String> template = new ReplyingKafkaTemplate<>(pf, repliesContainer);
    // 设置回复超时,和生产者端get()方法的超时保持一致
    template.setReplyTimeout(10000);
    return template;
}

4. 移除无效配置

你在kafkaListenerContainerFactory中设置的setReplyTemplate是给@KafkaListener注解使用的,自定义监听器不需要此配置,可移除避免混淆:

public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() throws IOException {
    ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());
    factory.setConcurrency(2);
    factory.setStatefulRetry(true);
    // 移除该行,自定义监听器不需要此配置
    // factory.setReplyTemplate(kafkaTemplate);
    factory.getContainerProperties().setAckMode(AckMode.MANUAL);
    return factory;
}

关键注意事项

  • 确保回复主题kReplies已存在,或Kafka开启了自动创建主题配置。
  • 必须保证响应消息携带的correlation-id和请求消息的完全一致,否则ReplyingKafkaTemplate无法关联请求和响应。
  • 回复容器的消费者组ID需唯一,避免和其他消费者组冲突。

内容的提问来源于stack exchange,提问作者user3672677

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 01:11:02