You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Boot中Kafka AdminClient单元测试空指针异常问题

Kafka TopicService单元测试Mock AdminClient时出现NullPointerException

问题场景

在为TopicService类编写单元测试,尝试Mock AdminClient调用createTopic方法时,触发了NullPointerException。

相关代码

TopicService类

@Service
public class TopicService {

  private static final Logger LOG = LoggerFactory.getLogger(TopicService.class);
  @Autowired
  private AdminClient adminClient;

  public void createTopic(Topic topic) throws ExecutionException, InterruptedException {
    adminClient
            .createTopics(Collections.singletonList(ServiceHelper.fromTopic(topic)))
            .values()
            .get(topic.getName())
            .get();
  }
}

单元测试类

package org.kafka.service;

import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.admin.CreateTopicsResult;
import org.apache.kafka.clients.admin.ListTopicsResult;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.kafka.common.KafkaFuture;
import org.apache.kafka.common.internals.KafkaFutureImpl;
import org.junit.Before;
import org.junit.BeforeClass;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.kafka.model.Topic;
import org.mockito.Mockito;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.boot.test.mock.mockito.MockBean;
import org.springframework.test.context.junit4.SpringRunner;

import java.util.*;
import java.util.concurrent.ExecutionException;

import static org.apache.kafka.clients.admin.AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG;
import static org.apache.kafka.clients.admin.AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG;
import static org.apache.kafka.common.internals.Topic.GROUP_METADATA_TOPIC_NAME;

@Slf4j
@RunWith(SpringRunner.class)
@SpringBootTest(classes = {TopicService.class})
public class TopicServiceTest {

    @Autowired
    TopicService topicService;

    @MockBean
    AdminClient adminClient;
    ListTopicsResult listTopicsResult;

    KafkaFuture<Set<String>> future;

    NewTopic newTopic;

    Topic topic;

    Collection<NewTopic> topicList;

    CreateTopicsResult createTopicsResult;
    Void t;

    Map<String,KafkaFuture<Void>> futureMap;

    private static final String TARGET_CONSUMER_GROUP_ID = "target-group-id";

    private static final Map<String, Object> CONF = new HashMap<>();

    @BeforeClass
    public static void createAdminClient() {
        try {
            CONF.put(BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
            CONF.put(REQUEST_TIMEOUT_MS_CONFIG, 120000);
            CONF.put("zookeeper.connect", "localhost:21891");

           AdminClient adminClient = AdminClient.create(CONF);
        } catch (Exception e) {
            throw new RuntimeException("create kafka admin client error", e);
        }
    }

    @Before
    public void setUp(){
        topicList = new ArrayList<>();
        newTopic = new NewTopic("topic-7",1, (short) 1);
        topicList.add(newTopic);
        futureMap = new HashMap<>();
        topic = new Topic();
        topic.setName("topic-1");
    }

    @Test
    public void createTopic() throws ExecutionException, InterruptedException {

        Properties consumerProperties = new Properties();

        Mockito.when(adminClient.createTopics(topicList))
                .thenReturn(Mockito.mock(CreateTopicsResult.class));

        Mockito.when(adminClient.createTopics(topicList).values())
                .thenReturn(Mockito.mock(Map.class));

        Mockito.when(adminClient.createTopics(topicList)
                .values()
                .get(GROUP_METADATA_TOPIC_NAME)).thenReturn(Mockito.mock(KafkaFutureImpl.class));

        Mockito.when(adminClient.createTopics(topicList)
                .values()
                .get(GROUP_METADATA_TOPIC_NAME)
                .get()).thenReturn(t);
         topicService.createTopic(topic);
    }
}

AdminClient配置类

package org.kafka.config;

import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.admin.AdminClientConfig;
import org.kafka.reader.Kafka;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Bean;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.kafka.core.KafkaAdmin;
import org.springframework.stereotype.Component;

import java.util.HashMap;
import java.util.Map;

@Component
public class AdminConfigurer {


    @Autowired
    private Kafka kafkaConfig;



    @Bean
    public Map<String, Object> kafkaAdminProperties() {
        final Map<String, Object> configs = new HashMap<>();
        configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaConfig.getBootstrapServers());
        if(kafkaConfig.getProperties().getSasl().getEnabled() && kafkaConfig.getSsl().getEnabled()) {
            configs.put("sasl.mechanism", kafkaConfig.getProperties().getSasl().getMechanism());
            configs.put("security.protocol", kafkaConfig.getProperties().getSasl().getSecurity().getProtocol());
            configs.put("ssl.keystore.location", kafkaConfig.getSsl().getKeystoreLocation());
            configs.put("ssl.keystore.password", kafkaConfig.getSsl().getKeystorePassword());
            configs.put("ssl.truststore.location", kafkaConfig.getSsl().getTruststoreLocation());
            configs.put("ssl.truststore.password", kafkaConfig.getSsl().getTruststorePassword());
            configs.put("sasl.jaas.config", String.format(kafkaConfig.getJaasTemplate(),
                    kafkaConfig.getProperties().getSasl().getJaas().getConfig().getUsername(),
                    kafkaConfig.getProperties().getSasl().getJaas().getConfig().getPassword()));
            configs.put("ssl.endpoint.identification.algorithm", "");
        }
        return configs;
    }


    @Bean
    public AdminClient getClient() {
        return AdminClient.create(kafkaAdminProperties());
    }

}

错误日志

java.lang.NullPointerException
    at org.kafka.service.TopicService.createTopic(TopicService.java:57)
    at org.kafka.service.TopicServiceTest.createTopic(TopicServiceTest.java:100)
    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.runners.model.FrameworkMethod$1.runReflectiveCall(FrameworkMethod.java:59)
    at org.junit.internal.runners.model.ReflectiveCallable.run(ReflectiveCallable.java:12)
    at org.junit.runners.model.FrameworkMethod.invokeExplosively(FrameworkMethod.java:56)
    at org.junit.internal.runners.statements.InvokeMethod.evaluate(InvokeMethod.java:17)

依赖信息

使用Spring版本2.7.1,相关依赖:

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>3.1.1</version>
    <scope>test</scope>
    <classifier>test</classifier>
</dependency>

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>

问题分析

  1. 参数不匹配:测试中Mock的是adminClient.createTopics(topicList),但实际TopicService调用的是createTopics(Collections.singletonList(ServiceHelper.fromTopic(topic))),传入的NewTopic列表和测试中定义的topicList不是同一个对象,Mockito默认按引用相等匹配参数,导致Mock方法未触发,返回null。
  2. 链式Mock错误:多次调用adminClient.createTopics(topicList)会生成不同的Mock实例,导致values()和后续get()调用无法关联到之前的Mock对象,最终返回null。
  3. Topic名称不匹配:测试中Mock的是get(GROUP_METADATA_TOPIC_NAME),但实际代码中是get(topic.getName())(即"topic-1"),两者不匹配,导致KafkaFuture为null,调用get()时触发NPE。

修复后的测试代码

@Slf4j
@RunWith(SpringRunner.class)
@SpringBootTest(classes = {TopicService.class})
public class TopicServiceTest {

    @Autowired
    TopicService topicService;

    @MockBean
    AdminClient adminClient;

    Topic topic;

    @Before
    public void setUp(){
        topic = new Topic();
        topic.setName("topic-1");
    }

    @Test
    public void createTopic() throws ExecutionException, InterruptedException {
        // 模拟ServiceHelper转换后的NewTopic
        NewTopic mockNewTopic = new NewTopic(topic.getName(), 1, (short)1);
        
        // 创建Mock的返回对象
        CreateTopicsResult mockCreateResult = Mockito.mock(CreateTopicsResult.class);
        KafkaFuture<Void> mockFuture = Mockito.mock(KafkaFuture.class);
        
        // 构建匹配topic名称的结果Map
        Map<String, KafkaFuture<Void>> resultMap = Collections.singletonMap(topic.getName(), mockFuture);
        
        // 按链式调用顺序Mock,使用参数匹配器兼容实际传入的集合
        Mockito.when(adminClient.createTopics(Mockito.anyCollection()))
                .thenReturn(mockCreateResult);
        Mockito.when(mockCreateResult.values()).thenReturn(resultMap);
        Mockito.when(mockFuture.get()).thenReturn(null); // Void类型返回null
        
        // 执行测试方法
        topicService.createTopic(topic);
        
        // 验证方法调用是否符合预期
        Mockito.verify(adminClient).createTopics(Mockito.singletonList(mockNewTopic));
        Mockito.verify(mockFuture).get();
    }
}

额外优化建议

  • 移除@BeforeClass中无用的AdminClient创建逻辑,测试中已通过@MockBean替换真实实例,该代码无意义。
  • 如果ServiceHelper.fromTopic逻辑复杂,可通过静态Mock隔离依赖:
    @MockStatic(ServiceHelper.class)
    public void createTopic() throws ExecutionException, InterruptedException {
        NewTopic mockNewTopic = new NewTopic(topic.getName(), 1, (short)1);
        Mockito.when(ServiceHelper.fromTopic(topic)).thenReturn(mockNewTopic);
        
        // 后续Mock逻辑同上
    }
    
  • 纯单元测试建议使用@ExtendWith(MockitoExtension.class)替代@RunWith(SpringRunner.class),无需加载Spring上下文,提升测试速度:
    @Slf4j
    @ExtendWith(MockitoExtension.class)
    public class TopicServiceTest {
    
        @InjectMocks
        TopicService topicService;
    
        @Mock
        AdminClient adminClient;
    
        // ... 测试逻辑 ...
    }
    

内容的提问来源于stack exchange,提问作者Mangesh Pawar

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 14:35:50