Mock ReactiveKafkaProducerTemplate的send方法返回SenderResult问题
如何Mock ReactiveKafkaProducerTemplate的send方法返回SenderResult
我正在尝试对ReactiveKafkaProducerTemplate的send方法进行Mock,初始代码如下:
@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,请问这一需求是否可以实现?如果可以的话,能否提供相关的示例参考?我花费了大量时间寻找解决方案但一直没有找到。
第一次尝试及报错
我按照建议尝试了如下实现:
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)))这一行遇到了如下空指针异常,异常信息没有说明哪个对象为null,我也没有发现存在空值的对象。
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)
第二次调试及报错
我已经通过提供的代码片段成功创建了Mock,以下是我需要测试的业务代码:
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(); }
当我从测试类中调用该方法时,已经确认使用的是Mock模板,但在.doOnSuccess(senderResult -> log.info("sent {} offset : {}", metrics, senderResult.recordMetadata().offset()))这一行抛出了java.lang.NullPointerException。异常没有提供空值相关的详细信息,我已经确认consumerRecord和metrics都不为空。
问题根因与解决
最终我找到了问题根源:Mock配置和实际调用不匹配,实际代码的send方法需要3个参数,但我之前的Mock配置只匹配了2个参数的send方法。我将代码更新为如下内容后问题解决:
when(reactiveKafkaProducerTemplate.send(Mockito.anyString(),Mockito.anyString(), Mockito.anyList())).thenReturn(Mono.just(result));
内容的提问来源于stack exchange,提问作者perplexedDev
相关产品推荐
相关产品推荐

