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

Kafka容器集成测试:自动注入的Consumer无法启动

问题:集成测试中Kafka消费者无法启动

本地搭建Kafka环境时,消费者可正常启动监听Topic并调用对应服务,但在集成测试中,无论使用@EmbeddedKafka还是Kafka Testcontainers,消费者都完全没有启动痕迹——Kafka日志中看不到消费者,服务方法也未被调用。已确认配置、Topic名称、Group ID及端口均设置正确。


集成测试类代码

@Testcontainers
@EnableAutoConfiguration(exclude = KafkaAutoConfiguration.class)
@ContextConfiguration(initializers = SqlContainerConfiguration.class, classes = {
    KafkaConsumerConfiguration.class,
    KafkaProducerConfiguration.class,
    KafkaListenerService.class
})
@ActiveProfiles("integration")
@SpringBootTest(classes = {TestOrchestratorApplication.class, KafkaConsumerConfiguration.class},
    webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
@Import({KafkaProducerConfiguration.class, KafkaConsumerConfiguration.class, KafkaListenerService.class})
@DirtiesContext
public class KafkaIntegrationTest {

  @Autowired
  private KafkaProducerConfiguration kafkaProducerConfiguration;

  @Autowired
  private KafkaConsumerConfiguration kafkaConsumerConfiguration;

  @Autowired
  private KafkaListenerService kafkaListenerService;

  // 部分Mock API和接口省略

  private static final String TEST_TOPIC = "test-topic";

  @Container
  private static KafkaContainer kafkaContainer = new KafkaContainer("5.5.0")
                                                     .withEnv("KAFKA_AUTO_CREATE_TOPICS_ENABLE", "true");

  @DynamicPropertySource
  static void kafkaProperties(DynamicPropertyRegistry registry) {
    List<String> env = new ArrayList<>();
    env.add("KAFKA_LISTENERS=" + kafkaContainer.getBootstrapServers());
    env.add("KAFKA_ADVERTISED_LISTENERS=" + "PLAINTEXT://localhost:9093");
    env.add("KAFKA_CFG_ADVERTISED_LISTENERS=" + kafkaContainer.getBootstrapServers());
    kafkaContainer.setEnv(env);
    kafkaContainer.setEnv(env);
    // 手动启动容器
    kafkaContainer.start();
    String bootstrapServers = kafkaContainer.getBootstrapServers();
    // 移除PLAINTEXT://前缀
    if (bootstrapServers.startsWith("PLAINTEXT://")) {
      bootstrapServers = bootstrapServers.substring("PLAINTEXT://".length());
    }

    // 重复配置bootstrap-servers
    String finalBootstrapServers = bootstrapServers;
    registry.add("spring.kafka.bootstrap-servers", () -> finalBootstrapServers);
    String finalBootstrapServers1 = bootstrapServers;
    registry.add("spring.kafka.bootstrap-servers", () -> finalBootstrapServers1);
  }

  @Test
  public void testKafkaIntegration() throws Exception {
    KafkaTemplate<String, UserEvent> kafkaTemplate = new KafkaTemplate<>(kafkaProducerConfiguration.producerFactory());

    // Mock方法设置省略
    when(methodToBeMocked)...

    // 发送消息到Topic
    ListenableFuture<?> sendResult = kafkaTemplate.send(TEST_TOPIC, aMessageToBeSend);

    // 等待消息发送完成
    sendResult.get(10, TimeUnit.SECONDS);

    // 等待监听器被调用
    latch.await(10, TimeUnit.SECONDS);
    System.out.println(kafkaContainer.getLogs());
    if (sendResult.isDone() && !sendResult.isCancelled()) {
      // 验证Mock方法被调用
      verify(methodToBeMocked)...;
    } else {
      fail("Failed to send message");
    }
  }
}

Kafka消费者配置类代码

@Configuration
public class KafkaConsumerConfiguration {

  private static final Logger logger = LoggerFactory.getLogger(KafkaConsumerConfiguration.class);

  @Value("${spring.kafka.bootstrap-servers}")
  private String bootstrapServers;

  @Value("${spring.kafka.group-id}")
  private String groupId;

  private final Map<String, Object> consumerConfig;

  public KafkaConsumerConfiguration(
      @Value("${spring.kafka.bootstrap-servers}") String bootstrapServers,
      @Value("${spring.kafka.group-id}") String groupId
  ) {
    this.bootstrapServers = bootstrapServers;
    this.groupId = groupId;

    consumerConfig = Map.of(
        ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers,
        ConsumerConfig.GROUP_ID_CONFIG, groupId,
        ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest",
        ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, keyDeserializer,
        ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, valueDeserializer
    );
  }

  public static final Class<StringDeserializer> keyDeserializer = StringDeserializer.class;
  public static final Class<JsonDeserializer> valueDeserializer = JsonDeserializer.class;

  @Bean
  public ConsumerFactory<String, UserEvent> consumerFactory() {
    logger.info("Creating Consumer Factory with Configuration: {}", consumerConfig);
    return new DefaultKafkaConsumerFactory<>(consumerConfig, new StringDeserializer(), new JsonDeserializer<>());
  }

  public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, UserEvent>> factory(
      ConsumerFactory<String, UserEvent> consumerFactory
  ) {
    ConcurrentKafkaListenerContainerFactory<String, UserEvent> factory =
        new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory);
    return factory;
  }

  public void setBootstrapServers(String bootstrapServers) {
    this.bootstrapServers = bootstrapServers;
  }

  public void setGroupId(String groupId) {
    this.groupId = groupId;
  }

  public Map<String, Object> getConsumerConfig() {
    return consumerConfig;
  }
}

问题分析与修复方案

1. 核心问题:KafkaListener容器工厂未注册

KafkaConsumerConfiguration中的factory方法没有添加@Bean注解,Spring无法识别这个容器工厂,导致@KafkaListener注解的方法无法创建对应的消费者容器。

修复:给factory方法添加@Bean注解

@Bean
public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, UserEvent>> factory(
    ConsumerFactory<String, UserEvent> consumerFactory
) {
  ConcurrentKafkaListenerContainerFactory<String, UserEvent> factory =
      new ConcurrentKafkaListenerContainerFactory<>();
  factory.setConsumerFactory(consumerFactory);
  return factory;
}

2. 测试类配置冗余与错误

  • 重复导入配置类:@SpringBootTest、@ContextConfiguration、@Import重复导入相同类,导致上下文混乱
  • 错误排除自动配置:@EnableAutoConfiguration(exclude = KafkaAutoConfiguration.class)会禁用Spring Kafka的自动配置逻辑,除非你完全手动管理所有Kafka组件,否则不建议这么做
  • DynamicPropertySource逻辑错误:
    • 手动设置Kafka容器环境变量会覆盖Testcontainers的默认正确配置,导致网络连接问题
    • 重复添加spring.kafka.bootstrap-servers到配置注册表,无意义
    • 手动调用kafkaContainer.start(),Testcontainers会自动管理容器生命周期
    • 移除PLAINTEXT://前缀是错误的,Spring Kafka需要完整的带协议的地址

修复后的测试类核心配置

@Testcontainers
@ActiveProfiles("integration")
@SpringBootTest(classes = TestOrchestratorApplication.class, webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
@DirtiesContext
public class KafkaIntegrationTest {

  @Container
  private static KafkaContainer kafkaContainer = new KafkaContainer("5.5.0")
                                                     .withEnv("KAFKA_AUTO_CREATE_TOPICS_ENABLE", "true");

  @DynamicPropertySource
  static void kafkaProperties(DynamicPropertyRegistry registry) {
    registry.add("spring.kafka.bootstrap-servers", kafkaContainer::getBootstrapServers);
    registry.add("spring.kafka.group-id", () -> "test-integration-group"); // 测试环境单独指定Group ID
  }

  // 其余测试代码省略
}

3. 其他潜在问题

  • Json反序列化配置:确保JsonDeserializer能正确处理UserEvent,可添加信任包配置:
    @Bean
    public ConsumerFactory<String, UserEvent> consumerFactory() {
      JsonDeserializer<UserEvent> deserializer = new JsonDeserializer<>(UserEvent.class);
      deserializer.addTrustedPackages("com.your.package.model"); // 替换为UserEvent所在包
      return new DefaultKafkaConsumerFactory<>(consumerConfig, new StringDeserializer(), deserializer);
    }
    
  • CountDownLatch使用:确保KafkaListenerService中的监听方法正确调用latch.countDown(),比如:
    @Service
    public class KafkaListenerService {
      private final CountDownLatch latch = new CountDownLatch(1);
    
      @KafkaListener(topics = "test-topic", groupId = "${spring.kafka.group-id}")
      public void handleUserEvent(UserEvent event) {
        // 业务逻辑
        latch.countDown();
      }
    
      public CountDownLatch getLatch() {
        return latch;
      }
    }
    
  • 测试环境配置冲突:检查application-integration.yml中是否有覆盖自定义Kafka配置的参数。

内容的提问来源于stack exchange,提问作者1029Coder

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 23:53:07