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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.21 14:57:57