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>
问题分析
- 参数不匹配:测试中Mock的是
adminClient.createTopics(topicList),但实际TopicService调用的是createTopics(Collections.singletonList(ServiceHelper.fromTopic(topic))),传入的NewTopic列表和测试中定义的topicList不是同一个对象,Mockito默认按引用相等匹配参数,导致Mock方法未触发,返回null。 - 链式Mock错误:多次调用
adminClient.createTopics(topicList)会生成不同的Mock实例,导致values()和后续get()调用无法关联到之前的Mock对象,最终返回null。 - 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
相关产品推荐
相关产品推荐

