无需@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
相关产品推荐
相关产品推荐

