【问题标题】:Simple embedded Kafka test example with spring boot带有 Spring Boot 的简单嵌入式 Kafka 测试示例
【发布时间】:2018-07-23 00:16:08
【问题描述】:

编辑仅供参考:working gitHub example


我在互联网上搜索,找不到嵌入式 Kafka 测试的有效且简单的示例。

我的设置是:

  • 弹簧靴
  • 多个@KafkaListener在一个班级中具有不同的主题
  • 用于测试的嵌入式 Kafka 开始正常
  • 使用发送到主题的 Kafka 模板进行测试,但 @KafkaListener 方法即使经过很长的睡眠时间也没有收到任何信息
  • 不显示警告或错误,日志中仅显示来自 Kafka 的垃圾信息

请帮助我。大多存在过度配置或过度设计的示例。我相信它可以简单地完成。 谢谢各位!

@Controller
public class KafkaController {

    private static final Logger LOG = getLogger(KafkaController.class);

    @KafkaListener(topics = "test.kafka.topic")
    public void receiveDunningHead(final String payload) {
        LOG.debug("Receiving event with payload [{}]", payload);
        //I will do database stuff here which i could check in db for testing
    }
}

私有静态字符串 SENDER_TOPIC = "test.kafka.topic";

@ClassRule
public static KafkaEmbedded embeddedKafka = new KafkaEmbedded(1, true, SENDER_TOPIC);

@Test
    public void testSend() throws InterruptedException, ExecutionException {

        Map<String, Object> senderProps = KafkaTestUtils.producerProps(embeddedKafka);

        KafkaProducer<Integer, String> producer = new KafkaProducer<>(senderProps);
        producer.send(new ProducerRecord<>(SENDER_TOPIC, 0, 0, "message00")).get();
        producer.send(new ProducerRecord<>(SENDER_TOPIC, 0, 1, "message01")).get();
        producer.send(new ProducerRecord<>(SENDER_TOPIC, 1, 0, "message10")).get();
        Thread.sleep(10000);
    }

【问题讨论】:

  • 显示代码。看看这是否有帮助stackoverflow.com/questions/48682745/…
  • @pvpkiran 这仍然不起作用。当我只是将发送部分带到我的主题时,测试只会测试自己,但永远不会到达我的 KafkaListener
  • 您的测试代码并不清楚KafkaController 是如何参与的。您如何确定侦听器已启动?
  • @ArtemBilan 因为方法上有 [@KafkaListener] 注释。还是我必须做点别的?
  • 对,测试需要使用该组件引导应用程序上下文

标签: spring-boot apache-kafka spring-kafka


【解决方案1】:

嵌入式 Kafka 测试适用于以下配置,

对测试类的注解

@EnableKafka
@SpringBootTest(classes = {KafkaController.class}) // Specify @KafkaListener class if its not the same class, or not loaded with test config
@EmbeddedKafka(
    partitions = 1, 
    controlledShutdown = false,
    brokerProperties = {
        "listeners=PLAINTEXT://localhost:3333", 
        "port=3333"
})
public class KafkaConsumerTest {
    @Autowired
    KafkaEmbedded kafkaEmbeded;

    @Autowired
    KafkaListenerEndpointRegistry kafkaListenerEndpointRegistry;

设置方法注释前

@Before
public void setUp() throws Exception {
  for (MessageListenerContainer messageListenerContainer : kafkaListenerEndpointRegistry.getListenerContainers()) {
    ContainerTestUtils.waitForAssignment(messageListenerContainer, 
    kafkaEmbeded.getPartitionsPerTopic());
  }
}

注意:我没有使用 @ClassRule 来创建嵌入式 Kafka,而是使用自动装配
@Autowired embeddedKafka

@Test
public void testReceive() throws Exception {
     kafkaTemplate.send(topic, data);
}

希望这会有所帮助!

编辑:测试用@TestConfiguration标记的配置类

@TestConfiguration
public class TestConfig {

@Bean
public ProducerFactory<String, String> producerFactory() {
    return new DefaultKafkaProducerFactory<>(KafkaTestUtils.producerProps(kafkaEmbedded));
}

@Bean
public KafkaTemplate<String, String> kafkaTemplate() {
    KafkaTemplate<String, String> kafkaTemplate = new KafkaTemplate<>(producerFactory());
    kafkaTemplate.setDefaultTopic(topic);
    return kafkaTemplate;
}

现在@Test 方法将自动装配 KafkaTemplate 并用于发送消息

kafkaTemplate.send(topic, data);

用上面的行更新了答案代码块

【讨论】:

  • 谢谢!这听起来不错,但是 [@EmbeddedKafka] 和 [kafkaListenerEndpointRegistry] 来自哪里?你能发布一个完整的导入示例吗?
  • 因为我们已经用@EnableKafka@EmbeddedKafka 注释我们的类,你可以在测试类中自动装配两者。在回答第一个代码块时,@Autowired KafkaEmbedded kafkaEmbedded 已经存在,就像您可以自动装配 kafkaListenerEndpointRegistry
  • 我仍然不知道“@EmbeddedKafka”来自哪里。为此需要哪个依赖项?我目前正在使用“spring-kafka-test”
  • TestConfig 最好在KafkaConsumerTest 中声明为内部类。在这种情况下: a) 必须是 static b) KafkaEmbedded 必须作为方法 producerFactory 的参数注入 c) 将 ProducerFactory 作为方法 kafkaTemplate 的参数注入,然后使用它代替打电话给producerFactory()
  • setup() 方法和ContainerTestUtils.waitForAssignment(..) 对我们来说是黄金。我们遇到了消费者在另一个测试类之后挂起,这导致消费者在下一个测试中没有收到任何东西。我们也使用@DirtiesContext(AFTER_CLASS)
【解决方案2】:

因为接受的答案不能编译或为我工作。我根据https://blog.mimacom.com/testing-apache-kafka-with-spring-boot/ 找到了另一个解决方案,我想与您分享。

依赖是'spring-kafka-test'版本:'2.2.7.RELEASE'

@RunWith(SpringRunner.class)
@EmbeddedKafka(partitions = 1, topics = { "testTopic" })
@SpringBootTest
public class SimpleKafkaTest {

    private static final String TEST_TOPIC = "testTopic";

    @Autowired
    EmbeddedKafkaBroker embeddedKafkaBroker;

    @Test
    public void testReceivingKafkaEvents() {
        Consumer<Integer, String> consumer = configureConsumer();
        Producer<Integer, String> producer = configureProducer();

        producer.send(new ProducerRecord<>(TEST_TOPIC, 123, "my-test-value"));

        ConsumerRecord<Integer, String> singleRecord = KafkaTestUtils.getSingleRecord(consumer, TEST_TOPIC);
        assertThat(singleRecord).isNotNull();
        assertThat(singleRecord.key()).isEqualTo(123);
        assertThat(singleRecord.value()).isEqualTo("my-test-value");

        consumer.close();
        producer.close();
    }

    private Consumer<Integer, String> configureConsumer() {
        Map<String, Object> consumerProps = KafkaTestUtils.consumerProps("testGroup", "true", embeddedKafkaBroker);
        consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        Consumer<Integer, String> consumer = new DefaultKafkaConsumerFactory<Integer, String>(consumerProps)
                .createConsumer();
        consumer.subscribe(Collections.singleton(TEST_TOPIC));
        return consumer;
    }

    private Producer<Integer, String> configureProducer() {
        Map<String, Object> producerProps = new HashMap<>(KafkaTestUtils.producerProps(embeddedKafkaBroker));
        return new DefaultKafkaProducerFactory<Integer, String>(producerProps).createProducer();
    }
}

【讨论】:

  • 你在测试嵌入式 kafka 本身吗?
  • java.lang.NoClassDefFoundError: kafka/common/KafkaException
【解决方案3】:

我现在解决了这个问题

@BeforeClass
public static void setUpBeforeClass() {
    System.setProperty("spring.kafka.bootstrap-servers", embeddedKafka.getBrokersAsString());
    System.setProperty("spring.cloud.stream.kafka.binder.zkNodes", embeddedKafka.getZookeeperConnectionString());
}

在调试时,我看到嵌入式 kaka 服务器占用了一个随机端口。

我找不到它的配置,所以我将 kafka 配置设置为与服务器相同。对我来说看起来还是有点难看。

我希望只有@Mayur 提到的那一行

@EmbeddedKafka(partitions = 1, controlledShutdown = false, brokerProperties = {"listeners=PLAINTEXT://localhost:9092", "port=9092"})

但在互联网上找不到正确的依赖项。

【讨论】:

  • 您可以在您的 application.properties 中设置 spring.kafka.bootstrap-servers=${spring.embedded.kafka.brokers} 进行测试,这应该可以。这是从 EmbeddedKafka 用它在启动时分配的随机端口填充的。
  • 这个注释@EmbeddedKafka在我的例子中来自spring-kafka-test-2.6.5。我在 pom 中依赖于 spring-kafka-test 并且我使用的是 spring-boot 2.4.2 版本。@Sylhare
【解决方案4】:

在集成测试中,不建议使用像 9092 这样的固定端口,因为多个测试应该能够灵活地从嵌入式实例打开自己的端口。所以,下面的实现是这样的,

注意:此实现基于 junit5(Jupiter:5.7.0) 和 spring-boot 2.3.4.RELEASE

测试类:

@EnableKafka
@SpringBootTest(classes = {ConsumerTest.Config.class, Consumer.class})
@EmbeddedKafka(
        partitions = 1,
        controlledShutdown = false)
@TestInstance(TestInstance.Lifecycle.PER_CLASS)
public class ConsumerTest {

    @Autowired
    private EmbeddedKafkaBroker kafkaEmbedded;

    @Autowired
    private KafkaListenerEndpointRegistry kafkaListenerEndpointRegistry;

    @BeforeAll
    public void setUp() throws Exception {
        for (final MessageListenerContainer messageListenerContainer : kafkaListenerEndpointRegistry.getListenerContainers()) {
            ContainerTestUtils.waitForAssignment(messageListenerContainer,
                    kafkaEmbedded.getPartitionsPerTopic());
        }
    }

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

    @Autowired
    private KafkaTemplate<String, Optional<Map<String, List<ImmutablePair<String, String>>>>> requestKafkaTemplate;

    @Test
    public void consume_success() {
        requestKafkaTemplate.send(topicName, load);
    }


    @Configuration
    @Import({
            KafkaListenerConfig.class,
            TopicConfig.class
    })
    public static class Config {

        @Value(value = "${spring.kafka.bootstrap-servers}")
        private String bootstrapAddress;

        @Bean
        public ProducerFactory<String, Optional<Map<String, List<ImmutablePair<String, String>>>>> requestProducerFactory() {
            final Map<String, Object> configProps = new HashMap<>();
            configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
            configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
            configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
            return new DefaultKafkaProducerFactory<>(configProps);
        }

        @Bean
        public KafkaTemplate<String, Optional<Map<String, List<ImmutablePair<String, String>>>>> requestKafkaTemplate() {
            return new KafkaTemplate<>(requestProducerFactory());
        }
    }
}

监听类:

@Component
public class Consumer {
    @KafkaListener(
            topics = "${topic.name}",
            containerFactory = "listenerContainerFactory"
    )
    @Override
    public void listener(
            final ConsumerRecord<String, Optional<Map<String, List<ImmutablePair<String, String>>>>> consumerRecord,
            final @Payload Optional<Map<String, List<ImmutablePair<String, String>>>> payload
    ) {
        
    }
}

监听器配置:

@Configuration
public class KafkaListenerConfig {

    @Value(value = "${spring.kafka.bootstrap-servers}")
    private String bootstrapAddress;

    @Value(value = "${topic.name}")
    private String resolvedTreeQueueName;

    @Bean
    public ConsumerFactory<String, Optional<Map<String, List<ImmutablePair<String, String>>>>> resolvedTreeConsumerFactory() {
        final Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, resolvedTreeQueueName);
        return new DefaultKafkaConsumerFactory<>(props, new StringDeserializer(), new CustomDeserializer());
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, Optional<Map<String, List<ImmutablePair<String, String>>>>> resolvedTreeListenerContainerFactory() {
        final ConcurrentKafkaListenerContainerFactory<String, Optional<Map<String, List<ImmutablePair<String, String>>>>> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(resolvedTreeConsumerFactory());
        return factory;
    }

}

主题配置:

@Configuration
public class TopicConfig {

    @Value(value = "${spring.kafka.bootstrap-servers}")
    private String bootstrapAddress;

    @Value(value = "${topic.name}")
    private String requestQueue;

    @Bean
    public KafkaAdmin kafkaAdmin() {
        Map<String, Object> configs = new HashMap<>();
        configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
        return new KafkaAdmin(configs);
    }

    @Bean
    public NewTopic requestTopic() {
        return new NewTopic(requestQueue, 1, (short) 1);
    }
}

application.properties:

spring.kafka.bootstrap-servers=${spring.embedded.kafka.brokers}

此分配是将嵌入式实例端口绑定到 KafkaTemplate 和 KafkaListners 的最重要的分配。

按照上面的实现,你可以为每个测试类打开动态端口,这样会更方便。

【讨论】:

    猜你喜欢
    • 2021-03-17
    • 1970-01-01
    • 2018-10-09
    • 1970-01-01
    • 2016-12-03
    • 2023-03-10
    • 1970-01-01
    • 2020-07-21
    • 2018-12-13
    相关资源
    最近更新 更多