代码之家  ›  专栏  ›  技术社区  ›  Sujit

春季卡夫卡消费者/生产者测试

  •  1
  • Sujit  · 技术社区  · 7 年前

    spring-kafka 卡夫卡传播的抽象。我能够整合制片人和;然而,从实际实现的角度来看,我不确定如何测试(特别是集成测试)消费者的业务逻辑 @KafkaListener . 我试着跟着 spring-kafk 关于这个主题的文档和各种博客,但没有一个能回答我的问题。

    Spring启动测试类

    //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 但是,我的问题是如何验证类中的业务逻辑是否使用 @卡夫卡 正确执行,并根据错误处理和其他业务场景将消息路由到不同的主题。在一些例子中,我看到 CountDownLatch Async 那么,如何断言执行,还不确定。

    任何帮助,谢谢。

    1 回复  |  直到 7 年前
        1
  •  2
  •   Gary Russell    7 年前

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

    集成测试可以使用该“不同”主题来断言侦听器按预期处理了它。

    您还可以添加一个 BeanPostProcessor 添加到您的测试用例并包装 ConsumerMessageListener

    编辑

    下面是一个在代理中包装侦听器的示例。。。

    @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();
                    }
                };
            }
    
        }
    
    }
    
    推荐文章