Spring Kafka报错a KafkaTemplate is required to support replies
问题根因
启动抛出java.lang.IllegalStateException: a KafkaTemplate is required to support replies是配置缺失+配置错误共同导致的:
- 没有显式声明和
ReplyingKafkaTemplate泛型匹配的KafkaTemplateBean,Spring上下文初始化时无法给回复相关组件注入所需的模板实例 - 回复容器配置错误:一是错误监听了业务生产主题
mytopic,会和正常业务消息消费逻辑冲突;二是没有给回复容器设置关联的KafkaTemplate,无法完成回复消息的转发、匹配流程 - 额外逻辑错误:请求-回复模式下的回复容器必须监听独立的专用回复主题,不能和业务发送主题共用,否则会出现消息被错误消费、请求和响应无法匹配的问题
修复方案
修改Kafka配置类,补全缺失的Bean声明,修正回复容器配置,参考代码如下:
@Configuration @EnableKafka public class KafkaConfig { // 显式声明通用KafkaTemplate Bean,泛型和ReplyingKafkaTemplate保持一致 @Bean public KafkaTemplate<Object, KafkaExampleRecord> kafkaTemplate(ProducerFactory<Object, KafkaExampleRecord> producerFactory) { return new KafkaTemplate<>(producerFactory); } @Bean public ReplyingKafkaTemplate<Object, KafkaExampleRecord, KafkaExampleRecord> replyingKafkaTemplate(ProducerFactory<Object, KafkaExampleRecord> producerFactory, ConcurrentMessageListenerContainer<Object, KafkaExampleRecord> repliesContainer) { ReplyingKafkaTemplate<Object, KafkaExampleRecord, KafkaExampleRecord> replyTemplate = new ReplyingKafkaTemplate<>(producerFactory, repliesContainer); // 可按需设置默认的回复超时时间 replyTemplate.setDefaultReplyTimeout(Duration.ofSeconds(10)); return replyTemplate; } @Bean public ConcurrentMessageListenerContainer<Object, KafkaExampleRecord> repliesContainer( ConcurrentKafkaListenerContainerFactory<Object, KafkaExampleRecord> containerFactory, KafkaTemplate<Object, KafkaExampleRecord> kafkaTemplate) { // 必须使用独立的回复主题,不要和业务主题mytopic共用,示例命名为mytopic-replies ConcurrentMessageListenerContainer<Object, KafkaExampleRecord> replyContainer = containerFactory.createContainer("mytopic-replies"); replyContainer.setAutoStartup(false); // 给回复容器绑定KafkaTemplate,满足回复消息处理要求 replyContainer.getContainerProperties().setKafkaTemplate(kafkaTemplate); return replyContainer; } }
后续配置注意事项
- 提前在Kafka集群中创建专用的回复主题
mytopic-replies,避免主题不存在导致容器启动失败 - 业务消费端处理完请求消息后,需要将响应结果发送到请求头中
KafkaHeaders.REPLY_TOPIC指定的主题,同时原样带回请求中的KafkaHeaders.CORRELATION_ID请求头,否则ReplyingKafkaTemplate无法将请求和对应的响应做匹配 - 现有控制器中的
requestReply方法逻辑不完整,调用sendAndReceive后除了获取发送结果,还需要调用replyFuture.get()拿到实际的响应内容做后续处理
内容的提问来源于stack exchange,提问作者Tony Stark
相关产品推荐
相关产品推荐

