Spring Kafka集成测试生产者发送超时(主题元数据不存在)问题排查与修复咨询
Spring Kafka集成测试生产者发送超时(主题元数据不存在)问题排查与修复咨询
嗨,我来帮你分析下这个问题~从你给出的代码和报错信息来看,核心问题是你的生产者服务在测试时连接的是本地真实Kafka(localhost:29092),而不是测试用的EmbeddedKafkaBroker,这就导致生产者找不到仅在EmbeddedKafka里创建的主题,最终触发超时错误。
问题原因拆解
- 生产者配置与测试环境不匹配:你在
KafkaProducerConfiguration里通过@Value读取了配置文件中的spring.kafka.producer.bootstrap.servers=localhost:29092,而测试类里的EmbeddedKafkaBroker是一个独立的嵌入式Kafka实例,有自己的专属连接地址。生产者连到了外部Kafka,自然找不到仅在EmbeddedKafka中存在的dpp_cil.dpp.event.external.downstream_response主题。 - 测试类的主题检查仅验证了EmbeddedKafka:
testTopicCreation()用EmbeddedKafkaBroker.getBrokersAsString()获取地址创建AdminClient,所以能查到主题;但实际发送的生产者用的是外部Kafka地址,两者不是同一个实例,这就是为什么一个测试通过另一个却失败。
修复方案
你需要让测试环境下的生产者连接到EmbeddedKafkaBroker的地址,最简洁的方式是用@DynamicPropertySource动态覆盖配置:
步骤1:修改测试类,添加动态配置
在ProducerServiceIntegrationTest类中新增动态配置方法,覆盖生产者的bootstrap servers配置:
import org.springframework.test.context.DynamicPropertyRegistry; import org.springframework.test.context.DynamicPropertySource; // ... 保留原有注解 public class ProducerServiceIntegrationTest { private static final String TOPIC_EXAMPLE_EXTERNE = "dpp_cil.dpp.event.external.downstream_response"; @Autowired private EmbeddedKafkaBroker embeddedKafkaBroker; @Autowired private DownstreamActionResponseProducer downstreamActionResponseProducer; // 新增:动态设置生产者连接到嵌入式Kafka @DynamicPropertySource static void kafkaProperties(DynamicPropertyRegistry registry) { registry.add("spring.kafka.producer.bootstrap.servers", embeddedKafkaBroker::getBrokersAsString); // 若测试需要,也可覆盖schema registry地址(当前用本地8081,根据实际情况调整) // registry.add("schema.registry.url", () -> "http://localhost:8081"); } // ... 保留原有测试方法 }
步骤2:简化主题创建(可选)
因为你已经用@EmbeddedKafka(topics = {"dpp_cil.dpp.event.external.downstream_response"})指定了要创建的主题,setUp()里的主题检查和创建逻辑其实是冗余的,可以直接删除——EmbeddedKafka会自动创建你指定的主题。
步骤3:消费者反序列化优化(可选)
你的消费者用了ErrorHandlingDeserializer,但需要确保它能正确处理Avro序列化的消息,建议补充委托反序列化器的配置:
import io.confluent.kafka.serializers.KafkaAvroDeserializer; import java.util.HashMap; import java.util.Map; // ... 原有消费者配置逻辑 Deserializer<DownstreamActionResponse> avroDeserializer = new KafkaAvroDeserializer(); Map<String, String> deserializerProps = new HashMap<>(); deserializerProps.put("schema.registry.url", "http://localhost:8081"); avroDeserializer.configure(deserializerProps, false); // 用ErrorHandlingDeserializer包装Avro反序列化器 Deserializer<DownstreamActionResponse> errorHandlingDeserializer = new ErrorHandlingDeserializer<>(avroDeserializer);
验证逻辑
修改后,生产者在测试时会连接到EmbeddedKafkaBroker的地址,和主题检查用的是同一个Kafka实例,发送消息时就能找到对应的主题,不会再触发超时错误。
备注:内容来源于stack exchange,提问作者Bhavana
相关产品推荐
相关产品推荐

