【问题标题】:Mock SenderResult in ReactiveKafkaProducerTemplate send method在 ReactiveKafkaProducerTemplate 发送方法中模拟 SenderResult
【发布时间】:2021-09-02 07:14:19
【问题描述】:

我正在尝试模拟响应式KafkaConsumerTemplate 的发送方法。

 @Mock
private ReactiveKafkaConsumerTemplate<String, String> reactiveKafkaConsumerTemplate;
@Mock
private ReactiveKafkaProducerTemplate<String, List<Object>> reactiveKafkaProducerTemplate;

 Mockito.when(reactiveKafkaConsumerTemplate.receiveAutoAck())
        .thenReturn(createConsumerRecords(2));



Mockito.when(reactiveKafkaProducerTemplate
.send(Mockito.anyString(),Mockito.anyString(),Mockito.anyList()))
                    .thenReturn(???);

我正在尝试模拟 reactiveProducerTemplate 的 send 方法以返回 SenderResult。有可能这样做吗?如果是的话,有人可以指点我的文档/样本来做到这一点。我花了很多时间寻找解决方案,但找不到任何解决方案。

更新:根据 Gary 的建议,我尝试了以下操作

ProducerRecord<String, List<Object>> record 
= new  ProducerRecord<String, List<Object>>(topic,"key", objectSetup.setup());
RecordMetadata meta 
= new RecordMetadata(new TopicPartition("topic",0),0,0,0,(long)1,2,1);
 
Mockito.when(reactiveKafkaProducerTemplate.send(topic,"key",objectSetup.setup())
.thenReturn(Mono.just(new SendResult<>(record, meta))));

我在 .thenReturn(Mono.just(new SendResult(record, meta)))) 行收到以下异常。它没有在异常中提到什么是空的,我也没有看到任何空的东西。

java.lang.NullPointerException
    at com.ServiceTests.cTestMethod(ServiceTests.java:69)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:498)
    at org.junit.platform.commons.util.ReflectionUtils.invokeMethod(ReflectionUtils.java:688)
    at org.junit.jupiter.engine.execution.MethodInvocation.proceed(MethodInvocation.java:60)
    at org.junit.jupiter.engine.execution.InvocationInterceptorChain$ValidatingInvocation.proceed(InvocationInterceptorChain.java:131)
    at org.junit.jupiter.engine.extension.TimeoutExtension.intercept(TimeoutExtension.java:149)
    at org.junit.jupiter.engine.extension.TimeoutExtension.interceptTestableMethod(TimeoutExtension.java:140)
    at org.junit.jupiter.engine.extension.TimeoutExtension.interceptTestMethod(TimeoutExtension.java:84)
    at org.junit.jupiter.engine.execution.ExecutableInvoker$ReflectiveInterceptorCall.lambda$ofVoidMethod$0(ExecutableInvoker.java:115)
    at org.junit.jupiter.engine.execution.ExecutableInvoker.lambda$invoke$0(ExecutableInvoker.java:105)
    at org.junit.jupiter.engine.execution.InvocationInterceptorChain$InterceptedInvocation.proceed(InvocationInterceptorChain.java:106)
    at org.junit.jupiter.engine.execution.InvocationInterceptorChain.proceed(InvocationInterceptorChain.java:64)
    at org.junit.jupiter.engine.execution.InvocationInterceptorChain.chainAndInvoke(InvocationInterceptorChain.java:45)
    at org.junit.jupiter.engine.execution.InvocationInterceptorChain.invoke(InvocationInterceptorChain.java:37)
    at org.junit.jupiter.engine.execution.ExecutableInvoker.invoke(ExecutableInvoker.java:104)
    at org.junit.jupiter.engine.execution.ExecutableInvoker.invoke(ExecutableInvoker.java:98)
    at org.junit.jupiter.engine.descriptor.TestMethodTestDescriptor.lambda$invokeTestMethod$6(TestMethodTestDescriptor.java:210)
    at org.junit.platform.engine.support.hierarchical.ThrowableCollector.execute(ThrowableCollector.java:73)
    at org.junit.jupiter.engine.descriptor.TestMethodTestDescriptor.invokeTestMethod(TestMethodTestDescriptor.java:206)
    at org.junit.jupiter.engine.descriptor.TestMethodTestDescriptor.execute(TestMethodTestDescriptor.java:131)
    at org.junit.jupiter.engine.descriptor.TestMethodTestDescriptor.execute(TestMethodTestDescriptor.java:65)
    at org.junit.platform.engine.support.hierarchical.NodeTestTask.lambda$executeRecursively$5(NodeTestTask.java:139)
    at org.junit.platform.engine.support.hierarchical.ThrowableCollector.execute(ThrowableCollector.java:73)
    at org.junit.platform.engine.support.hierarchical.NodeTestTask.lambda$executeRecursively$7(NodeTestTask.java:129)
    at org.junit.platform.engine.support.hierarchical.Node.around(Node.java:137)
    at org.junit.platform.engine.support.hierarchical.NodeTestTask.lambda$executeRecursively$8(NodeTestTask.java:127)
    at org.junit.platform.engine.support.hierarchical.ThrowableCollector.execute(ThrowableCollector.java:73)
    at org.junit.platform.engine.support.hierarchical.NodeTestTask.executeRecursively(NodeTestTask.java:126)
    at org.junit.platform.engine.support.hierarchical.NodeTestTask.execute(NodeTestTask.java:84)

更新 2:我可以使用 Gary 的代码 sn-p 创建模拟。这是我要测试的代码

  public void sendToKafka(ConsumerRecord<String, String> consumerRecord){
    log.info("sending to topic={}, {}={},", destinationTopic, Metric.class.getSimpleName(), consumerRecord);
    List<Object> metrics = transformRecord(consumerRecord);
    kafkaProducerTemplate.send(destinationTopic, consumerRecord.key(), metrics)
            .doOnSuccess(senderResult -> log.info("sent {} offset : {}", metrics, senderResult.recordMetadata().offset()))
            .doOnError(throwable -> log.error("Error while sending message to destination topic : {}", throwable.getMessage()))
            .subscribe();
}

当我从我的测试中调用这个方法时,我可以看到模板是模拟模板但是,我得到一个 java.lang.NullPointerException 就行了 .doOnSuccess(senderResult -> log.info("sent {} offset : {}", metrics, senderResult.recordMetadata().offset()))

这个异常没有给出关于什么是空的任何细节。我确认了 consumerRecord 和 metrics 不为空。

发现问题出在设置上。实际代码需要 3 个参数,并且在设置中我只模拟了 send 方法的 2 个参数。 将代码更新为:

when(reactiveKafkaProducerTemplate.send(Mockito.anyString(),Mockito.anyString(), Mockito.anyList())).thenReturn(Mono.just(result));

【问题讨论】:

    标签: apache-kafka spring-kafka project-reactor reactor-kafka


    【解决方案1】:
    @Test
    void test() {
        ReactiveKafkaProducerTemplate<String, String> template = mock(ReactiveKafkaProducerTemplate.class);
        RecordMetadata meta = new RecordMetadata(new TopicPartition("foo", 0), 0L, 0L, 0L, 0L, 0, 2);
        SenderResult result = mock(SenderResult.class);
        when(result.recordMetadata()).thenReturn(meta);
        when(template.send("foo", "bar")).thenReturn(Mono.just(result));
        template.send("foo", "bar")
                .doOnNext(sr -> {
                    assertThat(sr.recordMetadata().toString()).isEqualTo("foo-0@0");
                })
                .subscribe();
    }
    

    【讨论】:

    • 我试过这个,但我得到一个空指针异常,我无法弄清楚是什么原因造成的。我在问题中添加了详细的例外情况。
    • 对不起,我的答案中的班级名称有误;我已经添加了一个完整的测试。
    • 谢谢。那工作得很好。但是,我现在面临另一个问题。这个测试方法是测试一段写入 Kafka 的代码,当我从测试中调用该方法时,我得到 java.lang.NullPointerException。我在问题描述中的更新 2 下添加了详细信息。您对导致这种情况的原因有什么想法吗?
    • 我发现问题出在我的发送设置上。在实际代码中,它需要 3 个参数,而在设置中我只使用 2 个参数进行了模拟。 .when(reactiveKafkaProducerTemplate.send(Mockito.anyString(),Mockito.anyString(), Mockito.anyList())).thenReturn(Mono.just(result));
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-10-11
    • 1970-01-01
    相关资源
    最近更新 更多