【问题标题】:Need Help | UnitTest | Kafka custom ConcurrentKafkaListenerContainerFactory with Custom Record FilterStrategy需要帮助 |单元测试 | Kafka 自定义 ConcurrentKafkaListenerContainerFactory 和自定义记录 FilterStrategy
【发布时间】:2020-08-04 01:55:45
【问题描述】:

尝试使用自定义 RecordFilterStrategy 对 ConcurrentKafkaListenerContainerFactory 进行单元测试,但无法找到测试过滤器策略的最佳方法。

  class ConsumerConfiguration {

  @Bean
  public ConcurrentKafkaListenerContainerFactory<String, Message> messageListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory<String, Message> factory =
        new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());
    factory.setRecordFilterStrategy(consumerRecord -> consumerRecord.value().getType().contains("MYFILTER"));
    return factory;
  }

}

单元测试

  @Mock
  KafkaProperties kafkaProperties;
  ConsumerConfiguration configuration;
  @Test
  void testMessageListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory<String, Message> actual =
        configuration.messageListenerContainerFactory();
    assertEquals(kafkaProperties.buildConsumerProperties(),
        actual.getContainerProperties().getKafkaConsumerProperties());
  }

虽然这段代码对单元测试很有用,但是没有提供自定义 RecordFilterStrategy 的代码覆盖,即 lambda 函数,所以如果有人已经为 Kafka Listener Container Factory 完成了单元测试/代码覆盖,那么需要帮助,什么是最好的方法处理。

【问题讨论】:

    标签: java spring unit-testing mockito spring-kafka


    【解决方案1】:

    对于工厂配置的纯单元测试,过滤器在哪里

    record -> record.value().contains("foo")
    

    你可以做这样的事情......

    @SpringBootTest
    class So63239177ApplicationTests {
    
        @Autowired
        ConcurrentKafkaListenerContainerFactory<String, String> factory;
    
        private int called;
    
        @Test
        void testFilter() throws Exception {
            MethodKafkaListenerEndpoint<String, String> endpoint = new MethodKafkaListenerEndpoint<>();
            endpoint.setBean(this);
            endpoint.setMethod(getClass().getDeclaredMethod("listen", ConsumerRecord.class));
            endpoint.setMessageHandlerMethodFactory(new DefaultMessageHandlerMethodFactory());
            endpoint.setTopics("foo");
            ConcurrentMessageListenerContainer<String, String> container = factory.createListenerContainer(endpoint);
            @SuppressWarnings("unchecked")
            MessageListener<String, String> messageListener =
                    (MessageListener<String, String>) container.getContainerProperties().getMessageListener();
            messageListener.onMessage(new ConsumerRecord<>("foo", 0, 0L, null, "foo"));
            assertThat(this.called).isEqualTo(0);
            messageListener.onMessage(new ConsumerRecord<>("foo", 0, 0L, null, "bar"));
            assertThat(this.called).isEqualTo(1);
        }
    
        public void listen(ConsumerRecord<String, String> in) {
            this.called++;
        }
    
    }
    

    编辑

    这是一种稍微简单(也许更清洁)的方法......

    @SpringBootTest
    class So632391771ApplicationTests {
    
        @Autowired
        KafkaListenerAnnotationBeanPostProcessor<?, ?> bpp;
    
        @Autowired
        KafkaListenerEndpointRegistry registry;
    
        @Test
        void testFilter() throws Exception {
            TestListener listener = new TestListener();
            this.bpp.postProcessAfterInitialization(listener, "test");
            ContainerProperties containerProperties = registry.getListenerContainer("test")
                    .getContainerProperties();
            assertThat(containerProperties.getGroupId()).isEqualTo("test");
            @SuppressWarnings("unchecked")
            MessageListener<String, String> messageListener =
                        (MessageListener<String, String>) containerProperties.getMessageListener();
            messageListener.onMessage(new ConsumerRecord<>("foo", 0, 0L, null, "foo"));
            assertThat(listener.called).isEqualTo(0);
            messageListener.onMessage(new ConsumerRecord<>("foo", 0, 0L, null, "bar"));
            assertThat(listener.called).isEqualTo(1);
        }
    
    }
    
    class TestListener { // NOTE: Not a @Bean
    
        int called;
    
        @KafkaListener(id = "test", topics = "foo", autoStartup = "false")
        public void listen(String in) {
            this.called++;
        }
    
    }
    

    【讨论】:

    • 非常感谢 Gary,它帮了大忙!
    • 我添加了另一种方法;也许更容易。
    • @Garry,关于这个话题还有一个问题,不确定,如果可以在同一个线程中提问,我想寻找我的消费者到特定的偏移量,在 Spring 中处理它的最佳方法是什么-卡夫卡。
    • 您不应该在 cmets 中提出补充问题 - 它无助于人们找到问题和答案。如果您对此问题还有其他问题,请参阅the documentation,请提出新问题。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-07-02
    • 1970-01-01
    • 1970-01-01
    • 2010-09-16
    • 2012-02-16
    • 1970-01-01
    • 2014-08-15
    相关资源
    最近更新 更多