【问题标题】:Unable to unit test a @KafkaListener annotated method无法对 @KafkaListener 注释的方法进行单元测试
【发布时间】:2018-10-17 06:18:59
【问题描述】:

我正在尝试在 Spring 中对 kafka 消费者类进行单元测试。我想知道,如果将 kafka 消息发送到其主题,则正确调用了侦听器方法。我的消费者类是这样注释的:

@KafkaListener(topics = "${kafka.topics.myTopic}")
public void myKafkaMessageEvent(final String message) { ...

如果我 @Autowire 消费者,当我发送 kafka 消息时,侦听器方法被正确调用,但我不能断言该方法已被调用,因为该类不是模拟。

如果我模拟消费者,当我发送 kafka 消息时,根本不会调用侦听器方法。我可以直接调用该方法,并断言它有效,但这并不能满足我的要求,即在我向其主题发送 kafka 消息时检查该方法是否被调用。

现在我已经在消费者内部放置了一个计数器,并在每次调用侦听器方法时递增它,然后检查它的值是否已更改。为测试创建一个变量对我来说似乎是一个糟糕的解决方案。

有没有办法让被嘲笑的消费者也收到卡夫卡消息?或者以其他方式断言调用了非模拟消费者侦听器方法?

【问题讨论】:

    标签: spring unit-testing apache-kafka spring-kafka


    【解决方案1】:

    听起来您正在请求类似于我们在 Spring AMQP 测试框架中的内容:https://docs.spring.io/spring-amqp/docs/2.0.3.RELEASE/reference/html/_reference.html#test-harness

    所以,如果你不擅长使用额外的变量,你可以借用 solution 并实现你自己的“线束”。

    我认为这应该是对框架的一个很好的补充,所以请提出适当的issue,我们可以一起为公众带来这样的工具。

    更新

    所以,根据 Spring AMQP 基础,我在我的测试配置中这样做了:

    public static class KafkaListenerTestHarness extends KafkaListenerAnnotationBeanPostProcessor {
    
        private final Map<String, Object> listeners = new HashMap<>();
    
        @Override
        protected void processListener(MethodKafkaListenerEndpoint endpoint, KafkaListener kafkaListener,
                Object bean, Object adminTarget, String beanName) {
    
            bean = Mockito.spy(bean);
    
            this.listeners.put(kafkaListener.id(), bean);
    
            super.processListener(endpoint, kafkaListener, bean, adminTarget, beanName);
        }
    
        @SuppressWarnings("unchecked")
        public <T> T getSpy(String id) {
            return (T) this.listeners.get(id);
        }
    
    }
    
    ...
    
    @SuppressWarnings("rawtypes")
    @Bean(name = KafkaListenerConfigUtils.KAFKA_LISTENER_ANNOTATION_PROCESSOR_BEAN_NAME)
    @Role(BeanDefinition.ROLE_INFRASTRUCTURE)
    public static KafkaListenerTestHarness kafkaListenerAnnotationBeanPostProcessor() {
        return new KafkaListenerTestHarness();
    }
    

    然后在目标测​​试用例中我这样使用它:

    @Autowired
    private KafkaListenerTestHarness harness;
    ...
    Listener listener = this.harness.getSpy("foo");
    
    verify(listener, times(2)).listen1("foo");
    

    【讨论】:

    • 那里有KafkaListenerAnnotationBeanPostProcessor 来处理所有那些@KafkaListener 方法。您只需要遵循RabbitListenerTestHarness 逻辑即可。让我稍后在我的答案中展示一些简单的东西!
    • 请在我的回答中找到更新。
    • 对不起,我有点困惑......所以我有这个:@Autowired private TestConfig.KafkaListenerTestHarness harness; 和这个:@Autowired private ReceiveMessageRetrievedEventHandler receiver; 第二个是我的消费者使用监听器方法。我将如何在这里使用您的代码?还有foo 参数我也没有得到...感谢您的帮助,抱歉我遇到了一些困难...
    • 注意我如何将spy - 存储到KafkaListenerTestHarness .listeners Map 中。只有这样您才能访问spyfoo@KafkaListener 上的 id()。看看我在KafkaListenerTestHarness 中得到了什么映射键。
    • 我想我明白了...但是我现在无法摆脱 Another endpoint is already registered with id 错误...
    猜你喜欢
    • 2017-09-13
    • 2023-02-09
    • 2013-10-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-07-29
    • 2012-06-01
    • 1970-01-01
    相关资源
    最近更新 更多