Embedded Kafka测试IntelliJ运行正常 mvn test执行失败排查
Embedded Kafka 测试跨执行环境失败问题
我编写了若干Embedded Kafka测试用例,以下是其中最简单的一个:
@SpringBootTest @EmbeddedKafka(partitions = 1, brokerProperties = { "listeners=PLAINTEXT://localhost:9092", "port=9092" }) @DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_EACH_TEST_METHOD) public class KafkaConsumerTest { private static final Logger logger = LoggerFactory.getLogger(KafkaConsumerTest.class); @Autowired private KafkaTemplate<String, String> kafkaTemplate; @Autowired private KafkaListenerEndpointRegistry registry; @Autowired private MessageListener messageListener; @Value("${kafka.fulfillment-topic}") private String topic; @Test public void messageListenerShouldReceiveMessage() throws InterruptedException { var fulfillmentContainer = (ConcurrentMessageListenerContainer<?, ?>) registry.getListenerContainer("fulfillment-listener"); fulfillmentContainer.stop(); var latch = new CountDownLatch(1); fulfillmentContainer.getContainerProperties().setMessageListener( (AcknowledgingConsumerAwareMessageListener<String, String>) (consumerRecord, acknowledgment, consumer) -> latch.countDown()); fulfillmentContainer.start(); var result = kafkaTemplate.send(topic, "hello, world!"); Assertions.assertThat(latch.await(10, TimeUnit.MINUTES)).isTrue(); } }
此外,我在非测试类中配置了如下监听器:
@KafkaListener(id = "fulfillment-listener", topics = "${kafka.fulfillment-topic}", groupId = "test-group") public void listenForFulfillment(ConsumerRecord<?, String> consumerRecord, Acknowledgment acknowledgment) { acknowledgment.acknowledge(); taskExecutor.execute(() -> { // 业务逻辑处理 }); }
该测试的设计目的是拦截发往此监听器的调用,确认消息投递行为正常;如有需要我也可以将消息转发给原始监听器,相关逻辑已经在其他测试中实现。
目前遇到的核心问题:
- 通过IntelliJ运行测试时可正常通过,请求能被成功拦截
- 在终端运行该测试则会失败,消息会直接发送到原始监听器,完全不会经过我设置的新监听器
更异常的现象:如果我仅通过.stop()方法停止容器且不再重启,通过IntelliJ运行时消息会如预期无法到达原始监听器,但在终端执行mvn test时,即便容器已处于停止状态,消息仍会异常到达原始监听器。
我所有的Kafka测试都存在该问题,目前暂无解决方案:所有测试通过IntelliJ内置测试运行器执行时全部成功,通过mvn test执行时全部失败。
附加环境信息
JDK版本
➜ ~ java -version openjdk version "17.0.3" 2022-04-19 OpenJDK Runtime Environment (build 17.0.3+7-Debian-1deb11u1) OpenJDK 64-Bit Server VM (build 17.0.3+7-Debian-1deb11u1, mixed mode, sharing)
依赖配置
<parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>3.0.0-SNAPSHOT</version> <relativePath/> <!-- 从仓库查找父依赖 --> </parent> ... <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency>
内容的提问来源于stack exchange,提问作者Steven Hay
相关产品推荐
相关产品推荐

