mvn install运行测试时@Autowired注入Bean实例不一致问题
Spring Boot集成测试Maven执行时Kafka消费实例不一致问题
问题现象
编写基于JUnit的Spring Boot集成测试时,定义两个核心类:
KafkaTestHelper通用组件:标注@Component,内部维护ConcurrentHashMap<String, ConsumerRecord<String, String>>类型的私有final变量messages存储监听到的Kafka消息,通过@KafkaListener监听指定主题,消费到消息后存入messages集合,对外提供getMessageByTopic方法供查询指定主题的消费消息,代码如下:@Slf4j @Component public class KafkaTestHelper { private final ConcurrentHashMap<String, ConsumerRecord<String, String>> messages = new ConcurrentHashMap<>(); @KafkaListener(topics = "#{kafkaTestHelper.getTopics()}") public void onListen(ConsumerRecord<String, String> record) { log.info("#OnListen topic : {}, record : {}", record.topic(), record); this.messages.put(record.topic(), record); log.info("put on object {}", System.identityHashCode(messages)); printHashMap(); } public String getMessageByTopic(String topic) { log.info("#trying to consume from {}", topic); log.info("consuming from object {}", System.identityHashCode(messages)); printHashMap(); String message = getMessage(topic); log.info("#returning message : {}", message); return message; } }IntegrationTest集成测试类:使用@RunWith(SpringRunner.class)、@SpringBootTest注解启动Spring测试上下文,通过@Autowired注入KafkaTestHelper实例,测试方法中先向Kafka发送指定主题消息,再调用kafkaTestHelper.getMessageByTopic方法获取监听到的结果消息做非空断言,代码如下:@RunWith(SpringRunner.class) @SpringBootTest(classes = {ImsStockApplication.class,}, webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT) public class IntegrationTest { @Autowired private KafkaTestHelper kafkaTestHelper; @Test public void success() throws Exception { producer.send("someTopics").block(); String valueFail = kafkaTestHelper.getMessageByTopic("someTopic"); // mvn构建时始终为空 Assert.assertNotNull(valueFail); } }
测试逻辑为:测试类发送Kafka主题消息,业务逻辑消费该消息后产出结果消息发送到对应主题,结果消息由KafkaTestHelper的监听器监听存储,最终在测试方法中读取结果做断言校验。
运行差异:
- IDE中直接运行测试:监听器写入消息、测试方法读取消息操作的是同一个
messages实例,日志打印的System.identityHashCode(messages)值一致,测试可正常通过,运行日志如下:Receive from topic someTopics with message somerecords #OnListen topic : someTopics, record : somerecords put on object 1654595627 <- 相同对象ID #trying to consume someTopics consuming from object 1654595627 <- 相同对象ID - 执行
mvn install命令运行测试:监听器写入消息的messages实例和测试方法读取消息的messages实例的System.identityHashCode值不一致,导致读取到的消息始终为空,断言失败,运行日志如下:Receive from topic someTopics with message somerecords #OnListen topic someTopics, record : somerecords put on object 1869941502 <- 不同对象ID #trying to consume from someTopics consuming from object 1398482101 <- 不同对象ID
按照Spring Bean默认单例规则,@Autowired注入的应为同一个KafkaTestHelper实例,对应内部的messages对象也应唯一,该问题仅在执行mvn install时出现、IDE运行时正常。
根因分析
该问题本质是Maven执行全量测试时Kafka监听器线程残留,导致消息被旧上下文的Bean实例消费,具体原因如下:
- IDE运行测试时通常仅执行单个测试类,Spring测试上下文仅初始化一次,不存在上下文重复创建、衍生线程残留的问题,因此监听器持有的Bean实例和测试类注入的实例完全一致。
- Maven执行
install阶段会跑全量测试用例,默认配置下maven-surefire插件会复用fork的测试JVM,多个测试类可能因配置差异触发Spring上下文多次创建;而Kafka消费者监听器是异步守护线程,如果旧上下文关闭时没有正确终止监听器线程,该线程会持续持有旧KafkaTestHelper实例的引用消费消息。新启动的上下文中虽然会创建新的KafkaTestHelper实例注入到测试类,但消息已经被旧实例消费,因此新实例的messages集合始终为空,日志打印的对象哈希值自然不一致。 - 额外注意:代码中发送消息的主题为
someTopics(复数),读取时传入的主题为someTopic(单数),存在拼写不一致问题,但该问题不会导致对象实例不一致的现象。
修复方案
按以下优先级调整配置即可解决问题:
- 配置测试类执行完成后销毁上下文,避免线程残留
在集成测试类上添加@DirtiesContext注解,强制测试类执行完成后关闭Spring上下文、终止所有关联的Bean生命周期(包括Kafka监听器线程):@RunWith(SpringRunner.class) @SpringBootTest(classes = {ImsStockApplication.class,}, webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT) @DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS) // 新增注解 public class IntegrationTest { // 原有代码不变 } - 调整maven-surefire插件配置,避免测试进程复用导致的交叉污染
在项目pom.xml中配置surefire插件,关闭并行测试、禁止fork进程复用,每个测试类执行完成后直接销毁进程,从根源上杜绝线程残留:<plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-surefire-plugin</artifactId> <configuration> <forkCount>1</forkCount> <reuseForks>false</reuseForks> <parallel>none</parallel> </configuration> </plugin> - 添加异步消费等待逻辑,避免时序问题
Kafka消费是异步流程,发送消息后不要立刻读取结果,使用等待工具轮询直到消息被消费(推荐使用Awaitility),避免偶发的时序问题,同时统一主题名拼写:@Test public void success() throws Exception { producer.send("someTopic").block(); // 修正主题名拼写,和监听主题保持一致 // 最多等待10秒,直到查询到对应消息 Awaitility.await().atMost(10, TimeUnit.SECONDS) .untilAsserted(() -> Assert.assertNotNull(kafkaTestHelper.getMessageByTopic("someTopic"))); } - 校验Kafka监听器容器生命周期配置
确认项目中ConcurrentKafkaListenerContainerFactory配置没有关闭自动启动,且Spring上下文关闭时会自动调用监听器容器的stop()方法终止消费者线程,避免线程泄漏。
内容的提问来源于stack exchange,提问作者hphp
相关产品推荐
相关产品推荐

