【问题标题】:Get data from topic after pushing in @EmbeddedKafka in spring boot Junit在 Spring Boot Junit 中推入 @EmbeddedKafka 后从主题中获取数据
【发布时间】:2019-06-24 11:02:27
【问题描述】:

我正在为我的 Spring Boot 应用程序编写 Junit 测试用例(使用 @EmbeddedKafka),该应用程序广泛使用 Spring-kafka 与其他服务和其他操作进行通信。

一个典型的例子是从 kafka 中删除数据(我们在 kafka 中推送 null 消息)。

目前在 delete() 方法中,我们首先检查 kafka 中是否存在任何请求删除的消息。 然后我们在 Kafka 中为该消息键推送 null

为上述方法逻辑编写 Junit 的步骤。

@Test
public void test(){
   //Push a message to Kafka (id=1234)
   //call test method service.delete(1234);
       //internally service.delete(1234) checks/validate whether message exists in kafka and then push null to delete topic.
  //check delete topic for delete message received.
  // Assertions
}

这里的问题是 Kafka 总是抛出 message not found 异常。在 service.delete() 方法中。

在控制台中检查日志时。我发现我的生产者配置为 kafka 使用不同的端口,而消费者配置使用不同的端口。

我不确定我是否遗漏了一些细节,或者这种行为的原因是什么。 任何帮助将不胜感激。

【问题讨论】:

  • 您需要显示所有配置(和日志)。
  • 请看我的回答。

标签: spring-boot spring-kafka embedded-kafka spring-kafka-test


【解决方案1】:

我有一个简单的 Spring Boot 应用供您考虑:

@SpringBootApplication
public class SpringBootEmbeddedKafkaApplication {

    public static final String MY_TOPIC = "myTopic";

    public BlockingQueue<String> kafkaMessages = new LinkedBlockingQueue<>();

    public static void main(String[] args) {
        SpringApplication.run(SpringBootEmbeddedKafkaApplication.class, args);
    }

    @KafkaListener(topics = MY_TOPIC)
    public void listener(String payload) {
        this.kafkaMessages.add(payload);
    }

}

application.properties:

spring.kafka.consumer.group-id=myGroup
spring.kafka.consumer.auto-offset-reset=earliest

并测试:

@RunWith(SpringRunner.class)
@SpringBootTest(properties =
        "spring.kafka.bootstrapServers:${" + EmbeddedKafkaBroker.SPRING_EMBEDDED_KAFKA_BROKERS + "}")
@EmbeddedKafka(topics = SpringBootEmbeddedKafkaApplication.MY_TOPIC)
public class SpringBootEmbeddedKafkaApplicationTests {

    @Autowired
    private KafkaTemplate<Object, String> kafkaTemplate;

    @Autowired
    private SpringBootEmbeddedKafkaApplication kafkaApplication;

    @Test
    public void testListenerWithEmbeddedKafka() throws InterruptedException {
        String testMessage = "foo";
        this.kafkaTemplate.send(SpringBootEmbeddedKafkaApplication.MY_TOPIC, testMessage);

        assertThat(this.kafkaApplication.kafkaMessages.poll(10, TimeUnit.SECONDS)).isEqualTo(testMessage);
    }

}

注意spring.kafka.consumer.auto-offset-reset=earliest让消费者从分区的开头读取。

在测试中应用的另一个重要选项是:

@SpringBootTest(properties =
        "spring.kafka.bootstrapServers:${" + EmbeddedKafkaBroker.SPRING_EMBEDDED_KAFKA_BROKERS + "}")

@EmbeddedKafka 填充 spring.embedded.kafka.brokers 系统属性并使 Spring Boot 自动配置知道我们需要将其值复制到 spring.kafka.bootstrapServers 配置属性。

或根据我们的docs 提供其他选项:

static {
    System.setProperty(EmbeddedKafkaBroker.BROKER_LIST_PROPERTY, "spring.kafka.bootstrap-servers");
}

【讨论】:

  • 让我试试这个
猜你喜欢
  • 2020-01-19
  • 2020-08-17
  • 2022-07-29
  • 1970-01-01
  • 2020-04-19
  • 1970-01-01
  • 2019-02-07
  • 2020-07-23
  • 1970-01-01
相关资源
最近更新 更多