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

如何使用普通Mockito模拟Kafka AdminClient返回指定Topic

解决办法

前提说明

确认使用的Mockito版本≥3.4.0,该版本开始原生支持静态方法模拟,无需依赖PowerMockito即可适配Junit5。

方案一:重构代码提升可测性(推荐)

原方法的核心问题是直接在方法内部调用AdminClient.create()静态方法生成实例,导致依赖耦合无法直接模拟。通过抽离AdminClient创建逻辑为工厂类,即可轻松实现依赖注入和模拟:

步骤1:定义AdminClient工厂接口与实现

// 工厂接口
public interface AdminClientFactory {
    AdminClient create(String brokerAddresses);
}

// 默认实现类
public class DefaultAdminClientFactory implements AdminClientFactory {
    @Override
    public AdminClient create(String brokerAddresses) {
        Properties props = new Properties();
        props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, brokerAddresses);
        return AdminClient.create(props);
    }
}

步骤2:调整原业务类为依赖注入模式

public class TopicService {
    private final AdminClientFactory adminClientFactory;
    // 构造注入工厂实例
    public TopicService(AdminClientFactory adminClientFactory) {
        this.adminClientFactory = adminClientFactory;
    }

    public List<String> getTopicByAddresses(String brokerAddresses) throws TopicsNameNotFound {
        List<String> result = new ArrayList<>();
        // 替换原静态创建逻辑为工厂调用
        try (AdminClient adminClient = adminClientFactory.create(brokerAddresses)) {
            ListTopicsOptions listTopicsOptions = new ListTopicsOptions();
            listTopicsOptions.timeoutMs(5000);
            result = adminClient.listTopics(listTopicsOptions).names().get()
                    .stream().filter(StringUtils::hasLength).toList();

        } catch (InterruptedException ie) {
            Thread.currentThread().interrupt();
        } catch (Exception e) {
            log.error(e.getMessage());
            throw new TopicsNameNotFound("Topics name not found for connection name: " + brokerAddresses);
        }
        if (result.isEmpty()) {
            throw new TopicsNameNotFound("Topics name not found for connection name: " + brokerAddresses);
        }
        return result;
    }
}

步骤3:编写单元测试

import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.admin.ListTopicsResult;
import org.apache.kafka.common.KafkaFuture;
import java.util.Set;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.when;
import static org.junit.jupiter.api.Assertions.*;

@ExtendWith(MockitoExtension.class)
public class TopicServiceTest {
    @Mock
    private AdminClientFactory adminClientFactory;
    @Mock
    private AdminClient adminClient;
    @Mock
    private ListTopicsResult listTopicsResult;
    @InjectMocks
    private TopicService topicService;

    @Test
    public void testGetTopicByAddresses_ReturnsExpectedTopics() throws Exception {
        Set<String> expectedTopics = Set.of("topic1", "topic2", "topic3");
        KafkaFuture<Set<String>> kafkaFuture = KafkaFuture.completedFuture(expectedTopics);
        // 串接模拟逻辑
        when(adminClientFactory.create(anyString())).thenReturn(adminClient);
        when(adminClient.listTopics(any(ListTopicsOptions.class))).thenReturn(listTopicsResult);
        when(listTopicsResult.names()).thenReturn(kafkaFuture);

        List<String> result = topicService.getTopicByAddresses("test-broker:9092");
        // 结果断言
        assertEquals(3, result.size());
        assertTrue(result.containsAll(expectedTopics));
    }

    @Test
    public void testGetTopicByAddresses_EmptyTopics_ThrowsException() throws Exception {
        KafkaFuture<Set<String>> kafkaFuture = KafkaFuture.completedFuture(Set.of());
        when(adminClientFactory.create(anyString())).thenReturn(adminClient);
        when(adminClient.listTopics(any(ListTopicsOptions.class))).thenReturn(listTopicsResult);
        when(listTopicsResult.names()).thenReturn(kafkaFuture);

        // 异常断言
        assertThrows(TopicsNameNotFound.class, () -> topicService.getTopicByAddresses("test-broker:9092"));
    }
}

方案二:不改动原业务代码,直接模拟静态方法

如果暂时无法重构原代码,可以直接用Mockito的静态模拟能力实现测试:

import org.junit.jupiter.api.Test;
import org.mockito.MockedStatic;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.mockStatic;
import static org.junit.jupiter.api.Assertions.*;

public class TopicServiceTest {
    private final TopicService topicService = new TopicService();

    @Test
    public void testGetTopicByAddresses_ReturnsExpectedTopics() throws Exception {
        Set<String> expectedTopics = Set.of("topic1", "topic2");
        AdminClient mockAdminClient = Mockito.mock(AdminClient.class);
        ListTopicsResult mockListResult = Mockito.mock(ListTopicsResult.class);
        KafkaFuture<Set<String>> mockFuture = KafkaFuture.completedFuture(expectedTopics);

        when(mockAdminClient.listTopics(any(ListTopicsOptions.class))).thenReturn(mockListResult);
        when(mockListResult.names()).thenReturn(mockFuture);

        // 模拟AdminClient静态create方法
        try (MockedStatic<AdminClient> mockedAdminClient = mockStatic(AdminClient.class)) {
            mockedAdminClient.when(() -> AdminClient.create(any(Properties.class))).thenReturn(mockAdminClient);
            List<String> result = topicService.getTopicByAddresses("test:9092");
            assertEquals(2, result.size());
            assertTrue(result.containsAll(expectedTopics));
        }
    }
}

注意事项

  • 静态方法模拟必须包裹在try-with-resources块中,执行完成后自动释放模拟作用域,避免影响其他测试用例
  • 优先选择重构方案,代码可维护性和测试性能更优

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 06:06:04