【问题标题】:Can't use a @KafkaListener more than once in a Test不能在测试中多次使用 @KafkaListener
【发布时间】:2019-09-03 14:35:26
【问题描述】:

我们正在尝试测试一个 cloud-stream-kafka 应用程序,在测试中我们有多个发送消息的测试方法,以及一个等待响应的 @KafkaListener。

但是,第一个测试往往会通过,而第二个测试往往会失败。

任何指针将不胜感激。

@SpringBootTest(properties = "spring.kafka.bootstrap-servers=${spring.embedded.kafka.brokers}")
@EmbeddedKafka(topics = "input", partitions = 1)
@DirtiesContext
class EmbeddedKafkaListenerTest {

  private CountDownLatch latch;

  private String message;

  @BeforeEach
  void setUp() {
    this.message = null;
    this.latch = new CountDownLatch(1);
  }

  @Test
  void testSendFirstMessage(@Autowired KafkaTemplate<String, byte[]> template)
      throws InterruptedException {
    template.send("input", "Hello World 1".getBytes());
    assertTrue(latch.await(10, TimeUnit.SECONDS));
    assertEquals("Hello World 1", message);
  }

  @Test
  void testSendSecondMessage(@Autowired KafkaTemplate<String, byte[]> template)
      throws InterruptedException {
    template.send("input", "Hello World 2".getBytes());
    assertTrue(latch.await(10, TimeUnit.SECONDS));
    assertEquals("Hello World 2", message);
  }

  @KafkaListener(topics = "input", id = "kafka-listener-consumer")
  void listener(Message<byte[]> message) {
    this.message = new String(message.getPayload(), StandardCharsets.UTF_8);
    this.latch.countDown();
  }
}

似乎@KafkaListener 的实例正在为每个测试注册,因为我们注意到使用id 值会导致java.lang.IllegalStateException: Another endpoint is already registered with id 'kafka-listener-consumer'

当使用的消息传递框架是 RabbitMQ 时,我使用 @RabbitListener 进行了类似的测试。我希望我可以做类似的事情,因为一些测试用例涉及等待没有消息被发布,我们可以通过 assertFalse(latch.await(10, TimeUnit.SECONDS)) 来做到这一点

【问题讨论】:

    标签: spring-kafka spring-cloud-stream spring-kafka-test


    【解决方案1】:

    我认为你的@KafkaListener 方法应该进入@Configuration 类,否则EmbeddedKafkaListenerTest 确实是每个测试方法实例化的,因此@KafkaListener 被解析与你有测试方法一样多。

    另一种方法是使用@DirtiesContext(classMode = ClassMode.AFTER_EACH_TEST_METHOD),这样您不仅可以为每个测试方法获得新鲜的EmbeddedKafkaListenerTest 实例,还可以使用Spring ApplicationContext

    还有一种使用@TestInstance(TestInstance.Lifecycle.PER_CLASS) 的方法,因此整个测试套件只有一个EmbeddedKafkaListenerTest,而您的@KafkaListener 不会被多次解析。

    【讨论】:

    • 谢谢! @TestInstance 工作。我试过@DirtiesContext,但似乎并没有解决问题。但是,@TestInstance(Lifecycle.PER_CLASS) 做到了。
    • 对。最好有 per class:更好的性能,因为您不需要为每个方法创建类的实例。
    猜你喜欢
    • 2020-06-30
    • 2020-06-23
    • 1970-01-01
    • 1970-01-01
    • 2012-05-20
    • 1970-01-01
    • 2019-02-09
    • 2022-01-05
    • 2019-09-05
    相关资源
    最近更新 更多