【问题标题】:Spring Boot KafkaTemplate and KafkaListener test with EmbeddedKafka fails使用 EmbeddedKafka 进行 Spring Boot KafkaTemplate 和 KafkaListener 测试失败
【发布时间】:2022-01-05 12:38:23
【问题描述】:

我有 2 个 Spring Boot 应用程序,一个是 Kafka 发布者,另一个是消费者。我正在尝试编写集成测试以确保发送和接收事件。

在 IDE 中或从命令行运行时测试为绿色,而无需其他测试,例如 mvn test -Dtest=KafkaPublisherTest。但是,当我构建整个项目时,测试以org.awaitility.core.ConditionTimeoutException 失败。项目中有多个@EmbeddedKafka 测试。

在日志中的这些行之后测试卡住了:

2021-11-30 09:17:12.366  INFO 1437 --- [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer    : wages-consumer-test: partitions assigned: [wages-test-0, wages-test-1]
2021-11-30 09:17:14.464  INFO 1437 --- [er-event-thread] kafka.controller.KafkaController         : [Controller id=0] Processing automatic preferred replica leader election

如果您对如何测试这些东西有更好的想法,请分享。

这是测试的样子:

@SpringBootTest(properties = { "kafka.wages-topic.bootstrap-address=${spring.embedded.kafka.brokers}" })
@EmbeddedKafka(partitions = 1, topics = "${kafka.wages-topic.name}")
class KafkaPublisherTest {

    @Autowired
    private TestWageProcessor testWageProcessor;

    @Autowired
    private KafkaPublisher kafkaPublisher;

    @Autowired
    private KafkaTemplate<String, WageEvent> kafkaTemplate;

    @Test
    void publish() {
        Date date = new Date();
        WageCreateDto wageCreateDto = new WageCreateDto().setName("test").setSurname("test").setWage(BigDecimal.ONE).setEventTime(date);
        kafkaPublisher.publish(wageCreateDto);

        kafkaTemplate.flush();
        WageEvent expected = new WageEvent().setName("test").setSurname("test").setWage(BigDecimal.ONE).setEventTimeMillis(date.toInstant().toEpochMilli());

        await()
                .atLeast(Duration.ONE_HUNDRED_MILLISECONDS)
                .atMost(Duration.TEN_SECONDS)
                .with()
                .pollInterval(Duration.ONE_HUNDRED_MILLISECONDS)
                .until(testWageProcessor::getLastReceivedWageEvent, equalTo(expected));
    }
}

发布者配置:

@Configuration
@EnableConfigurationProperties(WagesTopicPublisherProperties.class)
public class KafkaConfiguration {

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

    @Bean
    public NewTopic wagesTopic(WagesTopicPublisherProperties wagesTopicPublisherProperties) {
        return new NewTopic(wagesTopicPublisherProperties.getName(), wagesTopicPublisherProperties.getPartitions(), wagesTopicPublisherProperties.getReplicationFactor());
    }

    @Primary
    @Bean
    public WageEventSerde wageEventSerde() {
        return new WageEventSerde();
    }

    @Bean
    public ProducerFactory<String, WageEvent> producerFactory(WagesTopicPublisherProperties wagesTopicPublisherProperties, WageEventSerde wageEventSerde) {
        Map<String, Object> configProps = new HashMap<>();
        configProps.put(
                ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
                wagesTopicPublisherProperties.getBootstrapAddress());
        configProps.put(
                ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
                StringSerializer.class);
        configProps.put(
                ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
                wageEventSerde.serializer().getClass());
        return new DefaultKafkaProducerFactory<>(configProps);
    }

    @Bean
    public KafkaTemplate<String, WageEvent> kafkaTemplate(ProducerFactory<String, WageEvent> producerFactory) {
        return new KafkaTemplate<>(producerFactory);
    }
}

消费者配置:

@Configuration
@EnableConfigurationProperties(WagesTopicConsumerProperties.class)
public class ConsumerConfiguration {

    @ConditionalOnMissingBean(WageEventSerde.class)
    @Bean
    public WageEventSerde wageEventSerde() {
        return new WageEventSerde();
    }

    @Bean
    public ConsumerFactory<String, WageEvent> wageConsumerFactory(WagesTopicConsumerProperties wagesTopicConsumerProperties, WageEventSerde wageEventSerde) {
        Map<String, Object> props = new HashMap<>();
        props.put(
                ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
                wagesTopicConsumerProperties.getBootstrapAddress());
        props.put(
                ConsumerConfig.GROUP_ID_CONFIG,
                wagesTopicConsumerProperties.getGroupId());
        props.put(
                ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
                StringDeserializer.class);
        props.put(
                ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
                wageEventSerde.deserializer().getClass());
        return new DefaultKafkaConsumerFactory<>(
                props,
                new StringDeserializer(),
                wageEventSerde.deserializer());
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, WageEvent> wageEventConcurrentKafkaListenerContainerFactory(ConsumerFactory<String, WageEvent> wageConsumerFactory) {

        ConcurrentKafkaListenerContainerFactory<String, WageEvent> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(wageConsumerFactory);
        return factory;
    }
}
    @KafkaListener(
            topics = "${kafka.wages-topic.name}",
            containerFactory = "wageEventConcurrentKafkaListenerContainerFactory")
    public void consumeWage(WageEvent wageEvent) {
        log.info("Wage event received: " + wageEvent);
        wageProcessor.process(wageEvent);
    }

这里是项目源码:https://github.com/aleksei17/springboot-rest-kafka-mysql

以下是构建失败的日志:https://drive.google.com/file/d/1uE2w8rmJhJy35s4UJXf4_ON3hs9JR6Au/view?usp=sharing

【问题讨论】:

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


【解决方案1】:

当我使用 Testcontainers Kafka 而不是 @EmbeddedKafka 时,问题就解决了。测试看起来像这样:

@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
class PublisherApplicationTest {

    public static final KafkaContainer kafka =
            new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka").withTag("5.4.3"));

    static {
        kafka.start();
        System.setProperty("kafka.wages-topic.bootstrap-address", kafka.getBootstrapServers());
    }

但是,我不能说我理解这个问题。当我使用here 描述的单例模式时,我遇到了同样的问题。也许像@DirtiesContext 这样的东西会有所帮助:它有助于解决工作中的一项测试,但在这个学习项目中没有。

【讨论】:

    猜你喜欢
    • 2020-11-12
    • 2020-06-23
    • 1970-01-01
    • 1970-01-01
    • 2021-05-02
    • 2019-04-01
    • 2018-08-29
    • 2021-09-12
    • 1970-01-01
    相关资源
    最近更新 更多