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

Mockito模拟Kafka Admin Client创建Topic时遇空指针问题及测试优化

关于Mockito模拟Kafka Admin Client创建Topic时的空指针问题及优化方案

问题场景

我正在使用Mockito编写单元测试,模拟Kafka Admin Client创建Topic的返回结果,但执行Mockito.when(mockCreateTopicResult.values()).thenReturn(mockKafkaTopicResult);时出现空指针异常,尽管已对mockKafkaTopicResult进行Mock。测试代码如下:

import java.util.ArrayList;
import java.util.Map;
import java.util.Set;
import java.util.TreeMap;
import java.util.concurrent.ExecutionException;

import org.apache.kafka.clients.admin.Admin;
import org.apache.kafka.clients.admin.CreateTopicsResult;
import org.apache.kafka.clients.admin.ListTopicsResult;
import org.apache.kafka.common.KafkaFuture;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.mockito.InjectMocks;
import org.mockito.Mockito;
import org.mockito.MockitoAnnotations;

import org.mockito.Mock;
import io.quarkus.test.junit.QuarkusTest;
import io.quarkus.test.junit.main.Launch;
import io.restassured.RestAssured;
import io.restassured.filter.log.RequestLoggingFilter;
import io.restassured.filter.log.ResponseLoggingFilter;

@QuarkusTest
public class AppTest {

    @InjectMocks
    private App mockApp;

    @Mock
    private Client mockClient;

    @Mock
    private Database mockDatabase;

    @Mock
    private Admin mockKafkaAdmin;

    @Mock
    private ListTopicsResult mockListTopicResult;

    @Mock
    private KafkaFuture<Set<String>> mockKafkaFuture;

    @Mock
    private KafkaFuture<Void> mockKafkaFutureVoid;

    @Mock
    private Set<String> mockSet;

    @Mock
    private CreateTopicsResult mockCreateTopicResult;

    @Mock
    private Map<String, KafkaFuture<Void>> mockKafkaTopicResult;
    
    @BeforeAll
    public static void setupAll() {
        RestAssured.filters(new RequestLoggingFilter(), new ResponseLoggingFilter());
    }

    @BeforeEach
    public void setUp() {
        MockitoAnnotations.openMocks(this);
    }

    @Test
    @Launch(value = {}, exitCode = 1)
    public void testLaunchCommandFailed() {}

    @Test
    public void testCreateTopic() throws InterruptedException, ExecutionException {
        Mockito.when(mockClient.getKafka()).thenReturn(mockKafkaAdmin);
        Mockito.when(mockClient.getKafka().createTopics(Mockito.anyList())).thenReturn(mockCreateTopicResult);
        Mockito.when(mockCreateTopicResult.values()).thenReturn(mockKafkaTopicResult);
        Mockito.when(mockKafkaTopicResult.get("meme.transmit.test")).thenReturn(mockKafkaFutureVoid);
        mockApp.createTopic("test");
        Mockito.when(mockClient.getKafka().listTopics()).thenReturn(mockListTopicResult);
        Mockito.when(mockClient.getKafka().listTopics().names()).thenReturn(mockKafkaFuture);
        Mockito.when(mockClient.getKafka().listTopics().names().get()).thenReturn(mockSet);
        Assertions.assertTrue(mockApp.containsTopic("test"));
    }
// ...
}

空指针异常原因分析

  1. Quarkus测试与Mockito注解冲突:在@QuarkusTest类中使用Mockito的@Mock和@InjectMocks时,Quarkus的测试容器依赖注入机制会和Mockito的注解初始化逻辑冲突,导致mockCreateTopicResult等mock实例未被正确初始化,调用values()方法时触发空指针。
  2. 链式调用冗余:代码中Mockito.when(mockClient.getKafka().createTopics(Mockito.anyList()))属于冗余链式调用,已mockmockClient.getKafka()返回mockKafkaAdmin,直接对mockKafkaAdmin的方法mock即可,重复调用可能引发潜在问题。

解决空指针问题的步骤

  1. 替换注解适配Quarkus测试机制:
    • 将Mockito的@Mock替换为Quarkus的@MockBean,确保mock实例被测试容器正确管理。
    • 移除@InjectMocks,改用@Inject注入被测试的App实例,Quarkus会自动将mock依赖注入目标类。
  2. 简化mock逻辑:直接对mockKafkaAdmin的方法进行mock,避免冗余链式调用。

修改后的核心测试代码示例:

@QuarkusTest
public class AppTest {

    @Inject
    private App app;

    @MockBean
    private Client client;

    @MockBean
    private Database database;

    @MockBean
    private Admin kafkaAdmin;

    @MockBean
    private CreateTopicsResult createTopicResult;

    @MockBean
    private Map<String, KafkaFuture<Void>> kafkaTopicResult;

    @MockBean
    private KafkaFuture<Void> mockKafkaFutureVoid;

    @MockBean
    private ListTopicsResult mockListTopicResult;

    @MockBean
    private KafkaFuture<Set<String>> mockKafkaFuture;

    @MockBean
    private Set<String> mockSet;
    
    @BeforeAll
    public static void setupAll() {
        RestAssured.filters(new RequestLoggingFilter(), new ResponseLoggingFilter());
    }

    @Test
    public void testCreateTopic() throws InterruptedException, ExecutionException {
        Mockito.when(client.getKafka()).thenReturn(kafkaAdmin);
        Mockito.when(kafkaAdmin.createTopics(Mockito.anyList())).thenReturn(createTopicResult);
        Mockito.when(createTopicResult.values()).thenReturn(kafkaTopicResult);
        Mockito.when(kafkaTopicResult.get("meme.transmit.test")).thenReturn(mockKafkaFutureVoid);

        app.createTopic("test");

        Mockito.when(kafkaAdmin.listTopics()).thenReturn(mockListTopicResult);
        Mockito.when(mockListTopicResult.names()).thenReturn(mockKafkaFuture);
        Mockito.when(mockKafkaFuture.get()).thenReturn(mockSet);
        Mockito.when(mockSet.contains("test")).thenReturn(true);

        Assertions.assertTrue(app.containsTopic("test"));
    }
}

更简便的Kafka Admin Client测试方案

1. 使用Testcontainers进行集成测试

如果需要贴近真实场景的测试,可通过Testcontainers启动嵌入式Kafka集群,直接调用真实Admin Client API,无需手动mock:

@Testcontainers
@QuarkusTest
public class AppKafkaIntegrationTest {

    @Container
    private static final KafkaContainer KAFKA = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:latest"));

    @Inject
    private App app;

    @BeforeAll
    static void setup() {
        System.setProperty("kafka.bootstrap.servers", KAFKA.getBootstrapServers());
    }

    @Test
    public void testCreateAndCheckTopic() throws InterruptedException, ExecutionException {
        app.createTopic("test");
        Assertions.assertTrue(app.containsTopic("test"));
    }
}

2. 简化Mockito mock逻辑

若坚持纯单元测试,可通过直接实现抽象类或使用KafkaFuture.completedFuture减少mock对象数量:

@Test
public void testCreateTopicSimplified() throws InterruptedException, ExecutionException {
    // 构建创建Topic的返回结果
    KafkaFuture<Void> completedFuture = KafkaFuture.completedFuture(null);
    Map<String, KafkaFuture<Void>> resultMap = Map.of("meme.transmit.test", completedFuture);

    Mockito.when(client.getKafka()).thenReturn(kafkaAdmin);
    Mockito.when(kafkaAdmin.createTopics(Mockito.anyList())).thenReturn(new CreateTopicsResult() {
        @Override
        public Map<String, KafkaFuture<Void>> values() {
            return resultMap;
        }
        @Override
        public KafkaFuture<Void> all() {
            return completedFuture;
        }
    });

    // 构建查询Topic的返回结果
    Set<String> topics = Set.of("test");
    KafkaFuture<Set<String>> listFuture = KafkaFuture.completedFuture(topics);
    ListTopicsResult listResult = Mockito.mock(ListTopicsResult.class);
    Mockito.when(listResult.names()).thenReturn(listFuture);
    Mockito.when(kafkaAdmin.listTopics()).thenReturn(listResult);

    app.createTopic("test");
    Assertions.assertTrue(app.containsTopic("test"));
}

内容的提问来源于stack exchange,提问作者Steven Garcia

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 07:10:32