如何使用普通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
相关产品推荐
相关产品推荐

