ReplyingKafkaTemplate触发KafkaReplyTimeoutException无返回问题排查
解决方案
- 补全生产者侧ReplyingKafkaTemplate的Bean声明
你当前配置类里的replyingKafkaTemplate方法没有添加@Bean注解,Spring不会将该实例托管到容器中,导致你的自定义配置完全不生效,添加注解即可:
- 补全生产者侧ReplyingKafkaTemplate的Bean声明
@Bean // 新增这个注解 public ReplyingKafkaTemplate<String, Object, Object> replyingKafkaTemplate(final ProducerFactory<String, Object> producerFactory, final ConcurrentKafkaListenerContainerFactory<String, Object> consumerFactory) { // 原有逻辑不变 }
- 检查消费者侧的注解配置
@KafkaHandler是类内方法的注解,你需要在承载该方法的类上添加@KafkaListener注解绑定监听主题,否则@SendTo注解不会生效:
- 检查消费者侧的注解配置
@KafkaListener(topics = "main-topic") // 类上新增这个注解 public class YourKafkaListener { // 原有@KafkaHandler方法不变 }
- 确保消费者侧存在可用的普通KafkaTemplate Bean
@SendTo自动发送回复依赖容器中存在默认的KafkaTemplate实例,如果消费者服务中没有配置该Bean,需要补充配置:
- 确保消费者侧存在可用的普通KafkaTemplate Bean
@Bean public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> producerFactory) { return new KafkaTemplate<>(producerFactory); }
同时确认该KafkaTemplate的value序列化器支持EntityVO类的序列化,否则发送回复时会抛出序列化异常导致消息无法发送。
- 调整@SendTo注解配置
如果使用的Spring Kafka版本较低,无参数的@SendTo可能无法正常读取请求头中的回复topic,可显式指定读取请求头的回复topic,或者直接固定回复topic测试连通性:
- 调整@SendTo注解配置
// 方式1:显式从请求头取回复topic @SendTo("#{requestHeaders['kafka_replyTopic']}") // 方式2:固定回复topic先测试连通性 // @SendTo("replies") @KafkaHandler public EntityVO[] queryAllEntity(final AllEntitiesQuery allEntitiesQuery, @Headers final Map<String, String> header) { // 原有逻辑不变 }
- 开启消费者侧Kafka操作日志排查问题
在消费者配置文件中添加日志配置,开启Kafka客户端的debug日志,可直观看到是否有发送回复消息的请求、以及发送失败的原因:
- 开启消费者侧Kafka操作日志排查问题
logging.level.org.springframework.kafka=debug logging.level.org.apache.kafka=debug
内容的提问来源于stack exchange,提问作者user1339768
相关产品推荐
相关产品推荐

