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

@SpringBootTest中指定containerFactory的@KafkaListener未触发问题

问题

使用@SpringBootTest结合Testcontainers为Kafka消费者编写集成测试时,通过KafkaTemplate发送的消息已成功写入主题,但配置了自定义containerFactory的@KafkaListener始终未触发消费。

相关代码如下:

Consumer类

@Component
@ConditionalOnProperty("outdoor.kafka.enableKafkaReading")
public class Consumer {

    @Value("${topic.name}")
    public final String topic;
    @Value("${consumer.group}")
    public final String consumerGroup;

    @KafkaListener(topics = "#{__listener.topic}", groupId = "#{__listener.consumerGroup}",
            containerFactory = "containerFactory")
    public void consume(String message) {
        LOGGER.info("Received message: {}", message);
    }
}

Kafka消费者配置类

@Configuration
public class KafkaConsumerConfig {
    // 实际使用自定义反序列化器,此处用String示例
    @Bean
    public ConsumerFactory<String, String> consumerFactory(ObjectMapper objectMapper) {
        return new DefaultKafkaConsumerFactory<>(
                kafkaProperties.buildConsumerProperties(null), new StringDeserializer(), new StringDeserializer());
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> containerFactory(CommonErrorHandler errorHandler,
                                                                                    ConsumerFactory<String, String> consumerFactory) {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(leaderboardStandingEventConsumerFactory);
        factory.setCommonErrorHandler(errorHandler);
        return factory;
    }
}

测试类

@SpringBootTest
@ContextConfiguration(initializers = TestContainersInitializer.class)
class ConsumerTest {
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    @SpyBean
    private Consumer consumer;

    @Test
    void consume() {
        Instant now = Instant.now();
        kafkaTemplate.send("topic-name", "1", "message");
        await()
            .pollInterval(Duration.ofSeconds(1))
            .atMost(5, SECONDS)
            .untilAsserted(() -> {
                verify(consumer).consume(any());
            });
   }
}

已确认消息写入主题,但监听器未启动消费,简化配置后问题仍存在。

原因分析

  • 配置类Bean引用错误:containerFactory方法中设置消费者工厂时,使用了未定义的leaderboardStandingEventConsumerFactory,而非注入的consumerFactory,导致容器工厂无法正确初始化消费者。
  • ConditionalOnProperty开关未开启:Consumer类上的@ConditionalOnProperty要求测试环境中outdoor.kafka.enableKafkaReading为true,否则Consumer不会被容器加载,监听器失效。
  • final字段注入顺序问题:直接给final字段加@Value注入,可能导致初始化顺序错误,@KafkaListener中的SpEL表达式无法正确获取主题和groupId值,监听器绑定的主题不匹配。
  • 测试主题硬编码不匹配:测试代码中硬编码发送到"topic-name",若配置文件中topic.name不是该值,会导致监听器监听的主题与发送主题不一致。

解决方案

  1. 修正配置类Bean引用
    将containerFactory方法中的错误Bean引用替换为注入的consumerFactory:
@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> containerFactory(CommonErrorHandler errorHandler,
                                                                                ConsumerFactory<String, String> consumerFactory) {
    ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory);
    factory.setCommonErrorHandler(errorHandler);
    return factory;
}
  1. 开启测试环境的Kafka消费开关
    在测试环境配置文件(如application-test.properties)中添加:
outdoor.kafka.enableKafkaReading=true
  1. 改用构造器注入final字段
    避免@Value直接注入final字段的初始化问题:
@Component
@ConditionalOnProperty("outdoor.kafka.enableKafkaReading")
public class Consumer {

    public final String topic;
    public final String consumerGroup;

    public Consumer(@Value("${topic.name}") String topic,
                   @Value("${consumer.group}") String consumerGroup) {
        this.topic = topic;
        this.consumerGroup = consumerGroup;
    }

    @KafkaListener(topics = "#{__listener.topic}", groupId = "#{__listener.consumerGroup}",
            containerFactory = "containerFactory")
    public void consume(String message) {
        LOGGER.info("Received message: {}", message);
    }
}
  1. 统一测试主题与配置
    测试中使用配置文件中的主题名称,避免硬编码:
@SpringBootTest
@ContextConfiguration(initializers = TestContainersInitializer.class)
class ConsumerTest {
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;

    @SpyBean
    private Consumer consumer;

    @Value("${topic.name}")
    private String testTopic;

    @Test
    void consume() {
        Instant now = Instant.now();
        kafkaTemplate.send(testTopic, "1", "message");
        await()
            .pollInterval(Duration.ofSeconds(1))
            .atMost(5, SECONDS)
            .untilAsserted(() -> {
                verify(consumer).consume(any());
            });
   }
}

内容的提问来源于stack exchange,提问作者Daniel Pop

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 16:05:09