【问题标题】:Kafka Consumer/Producer test in Spring KafkaSpring Kafka 中的 Kafka 消费者/生产者测试
【发布时间】:2019-05-09 18:24:53
【问题描述】:

我目前正在开发 Kafka 模块,我正在使用 spring-kafka Kafka 通信的抽象。我能够从实际实现的角度集成生产者和消费者,但是,我不确定如何使用@KafkaListener 测试(特别是集成测试)消费者周围的业务逻辑。我尝试关注spring-kafk 文档和有关该主题的各种博客,但没有一个回答我的预期问题。

Spring Boot 测试类

//imports not mentioned due to brevity

@RunWith(SpringRunner.class)
@SpringBootTest(classes = PaymentAccountUpdaterApplication.class,
                webEnvironment = SpringBootTest.WebEnvironment.NONE)
public class CardUpdaterMessagingIntegrationTest {

    private final static String cardUpdateTopic = "TP.PRF.CARDEVENTS";

    @Autowired
    private ObjectMapper objectMapper;

    @ClassRule
    public static KafkaEmbedded kafkaEmbedded =
            new KafkaEmbedded(1, false, cardUpdateTopic);

    @Test
    public void sampleTest() throws Exception {
        Map<String, Object> consumerConfig =
                KafkaTestUtils.consumerProps("test", "false", kafkaEmbedded);
        consumerConfig.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        consumerConfig.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);

        ConsumerFactory<String, String> cf = new DefaultKafkaConsumerFactory<>(consumerConfig);
        ContainerProperties containerProperties = new ContainerProperties(cardUpdateTopic);
        containerProperties.setMessageListener(new SafeStringJsonMessageConverter());
        KafkaMessageListenerContainer<String, String>
                container = new KafkaMessageListenerContainer<>(cf, containerProperties);

        BlockingQueue<ConsumerRecord<String, String>> records = new LinkedBlockingQueue<>();
        container.setupMessageListener((MessageListener<String, String>) data -> {
            System.out.println("Added to Queue: "+ data);
            records.add(data);
        });
        container.setBeanName("templateTests");
        container.start();
        ContainerTestUtils.waitForAssignment(container, kafkaEmbedded.getPartitionsPerTopic());


        Map<String, Object> producerConfig = KafkaTestUtils.senderProps(kafkaEmbedded.getBrokersAsString());
        producerConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        producerConfig.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);

        ProducerFactory<String, Object> pf =
                new DefaultKafkaProducerFactory<>(producerConfig);
        KafkaTemplate<String, Object> kafkaTemplate = new KafkaTemplate<>(pf);

        String payload = objectMapper.writeValueAsString(accountWrapper());
        kafkaTemplate.send(cardUpdateTopic, 0, payload);
        ConsumerRecord<String, String> received = records.poll(10, TimeUnit.SECONDS);

        assertThat(received).has(partition(0));
    }


    @After
    public void after() {
        kafkaEmbedded.after();
    }

    private AccountWrapper accountWrapper() {
        return AccountWrapper.builder()
                .eventSource("PROFILE")
                .eventName("INITIAL_LOAD_CARD")
                .eventTime(LocalDateTime.now().toString())
                .eventID("8730c547-02bd-45c0-857b-d90f859e886c")
                .details(AccountDetail.builder()
                        .customerId("idArZ_K2IgE86DcPhv-uZw")
                        .vaultId("912A60928AD04F69F3877D5B422327EE")
                        .expiryDate("122019")
                        .build())
                .build();
    }
}

监听类

@Service
public class ConsumerMessageListener {
    private static final Logger LOGGER = LoggerFactory.getLogger(ConsumerMessageListener.class);

    private ConsumerMessageProcessorService consumerMessageProcessorService;

    public ConsumerMessageListener(ConsumerMessageProcessorService consumerMessageProcessorService) {
        this.consumerMessageProcessorService = consumerMessageProcessorService;
    }


    @KafkaListener(id = "cardUpdateEventListener",
            topics = "${kafka.consumer.cardupdates.topic}",
            containerFactory = "kafkaJsonListenerContainerFactory")
    public void processIncomingMessage(Payload<AccountWrapper,Object> payloadContainer,
                                       Acknowledgment acknowledgment,
                                       @Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
                                       @Header(KafkaHeaders.RECEIVED_PARTITION_ID) String partitionId,
                                       @Header(KafkaHeaders.OFFSET) String offset) {

        try {
            // business logic to process the message
            consumerMessageProcessorService.processIncomingMessage(payloadContainer);
        } catch (Exception e) {
            LOGGER.error("Unhandled exception in card event message consumer. Discarding offset commit." +
                    "message:: {}, details:: {}", e.getMessage(), messageMetadataInfo);
            throw e;
        }
        acknowledgment.acknowledge();
    }
}

我的问题是:在测试类中,我正在断言从BlockingQueue 轮询的分区、有效负载等,但是,我的问题是如何验证我在使用@KafkaListener 注释的类中的业务逻辑正在获取正确执行并根据错误处理和其他业务场景将消息路由到不同的主题。在一些示例中,我看到CountDownLatch 断言我不想将其放入我的业务逻辑中以在生产级代码中断言。消息处理器也是Async 所以,如何断言执行,不确定。

任何帮助,不胜感激。

【问题讨论】:

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


    【解决方案1】:

    正在正确执行并根据错误处理和其他业务场景将消息路由到不同的主题。

    集成测试可以从那个“不同的”主题中消费,以断言侦听器按预期处理了它。

    您还可以将BeanPostProcessor 添加到您的测试用例中,并将ConsumerMessageListener bean 包装在代理中以验证输入参数是否符合预期。

    编辑

    这是一个将监听器包装在代理中的示例...

    @SpringBootApplication
    public class So53678801Application {
    
        public static void main(String[] args) {
            SpringApplication.run(So53678801Application.class, args);
        }
    
        @Bean
        public MessageConverter converter() {
            return new StringJsonMessageConverter();
        }
    
        public static class Foo {
    
            private String bar;
    
            public Foo() {
                super();
            }
    
            public Foo(String bar) {
                this.bar = bar;
            }
    
            public String getBar() {
                return this.bar;
            }
    
            public void setBar(String bar) {
                this.bar = bar;
            }
    
            @Override
            public String toString() {
                return "Foo [bar=" + this.bar + "]";
            }
    
        }
    
    }
    
    @Component
    class Listener {
    
        @KafkaListener(id = "so53678801", topics = "so53678801")
        public void processIncomingMessage(Foo payload,
                Acknowledgment acknowledgment,
                @Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
                @Header(KafkaHeaders.RECEIVED_PARTITION_ID) String partitionId,
                @Header(KafkaHeaders.OFFSET) String offset) {
    
            System.out.println(payload);
            // ...
            acknowledgment.acknowledge();
        }
    
    }
    

    spring.kafka.consumer.enable-auto-commit=false
    spring.kafka.consumer.auto-offset-reset=earliest
    spring.kafka.listener.ack-mode=manual
    

    @RunWith(SpringRunner.class)
    @SpringBootTest(classes = { So53678801Application.class,
            So53678801ApplicationTests.TestConfig.class})
    public class So53678801ApplicationTests {
    
        @ClassRule
        public static EmbeddedKafkaRule embededKafka = new EmbeddedKafkaRule(1, false, "so53678801");
    
        @BeforeClass
        public static void setup() {
            System.setProperty("spring.kafka.bootstrap-servers",
                    embededKafka.getEmbeddedKafka().getBrokersAsString());
        }
    
        @Autowired
        private KafkaTemplate<String, String> template;
    
        @Autowired
        private ListenerWrapper wrapper;
    
        @Test
        public void test() throws Exception {
            this.template.send("so53678801", "{\"bar\":\"baz\"}");
            assertThat(this.wrapper.latch.await(10, TimeUnit.SECONDS)).isTrue();
            assertThat(this.wrapper.argsReceived[0]).isInstanceOf(Foo.class);
            assertThat(((Foo) this.wrapper.argsReceived[0]).getBar()).isEqualTo("baz");
            assertThat(this.wrapper.ackCalled).isTrue();
        }
    
        @Configuration
        public static class TestConfig {
    
            @Bean
            public static ListenerWrapper bpp() { // BPPs have to be static
                return new ListenerWrapper();
            }
    
        }
    
        public static class ListenerWrapper implements BeanPostProcessor, Ordered {
    
            private final CountDownLatch latch = new CountDownLatch(1);
    
            private Object[] argsReceived;
    
            private boolean ackCalled;
    
            @Override
            public int getOrder() {
                return Ordered.HIGHEST_PRECEDENCE;
            }
    
            @Override
            public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException {
                if (bean instanceof Listener) {
                    ProxyFactory pf = new ProxyFactory(bean);
                    pf.setProxyTargetClass(true); // unless the listener is on an interface
                    pf.addAdvice(interceptor());
                    return pf.getProxy();
                }
                return bean;
            }
    
            private MethodInterceptor interceptor() {
                return invocation -> {
                    if (invocation.getMethod().getName().equals("processIncomingMessage")) {
                        Object[] args = invocation.getArguments();
                        this.argsReceived = Arrays.copyOf(args, args.length);
                        Acknowledgment ack = (Acknowledgment) args[1];
                        args[1] = (Acknowledgment) () -> {
                            this.ackCalled = true;
                            ack.acknowledge();
                        };
                        try {
                            return invocation.proceed();
                        }
                        finally {
                            this.latch.countDown();
                        }
                    }
                    else {
                        return invocation.proceed();
                    }
                };
            }
    
        }
    
    }
    

    【讨论】:

    • 谢谢加里。我将我的错误主题添加到嵌入式 Kafka 以在发送类似于事件主题时进行初始化和侦听器调用。这工作正常。接收到的对象是一个 Json,我计划根据我的场景对业务方法发送的各种属性进行解析和断言。不过我没看懂,通过 ConsumerMessageListener 验证输入参数,请您详细说明一下。
    • 查看我的答案的编辑,了解如何包装监听器的示例。
    猜你喜欢
    • 1970-01-01
    • 2014-04-09
    • 2018-12-18
    • 1970-01-01
    • 1970-01-01
    • 2018-04-28
    • 2015-03-25
    • 2020-05-21
    • 2017-11-03
    相关资源
    最近更新 更多