WebFlux中publishOn、subscribeOn及反应式多线程场景测试方案咨询
单元测试实现方案
首先要做依赖隔离,避免外部依赖影响测试结果:
- 将硬编码的
ASYNC_SCHEDULER改成可注入的依赖,不要用static final固定写死在业务代码里,单元测试时直接传入Schedulers.immediate(),让所有操作都同步运行在当前测试线程,不需要额外处理异步等待逻辑。 - 用Mock工具模拟外部HTTP调用,比如Spring自带的
MockRestClient,构造预设的HTTP响应,不需要真实发起网络请求。 - 业务方法要返回反应式类型(比如
Mono<StatusResponse>),用StepVerifier订阅返回的流来验证结果,不需要手动处理订阅逻辑。
单元测试示例代码:
// 改造后的业务代码,依赖全部可注入 public class StatusService { private final RestClient restClient; private final Scheduler asyncScheduler; private final String baseUrl; public StatusService(RestClient restClient, Scheduler asyncScheduler, @Value("${api.base-url}") String baseUrl) { this.restClient = restClient; this.asyncScheduler = asyncScheduler; this.baseUrl = baseUrl; } public Mono<ProcessedResult> getProcessedStatus() { return restClient .get() .uri(baseUrl + "/id") .retrieve() .bodyToMono(StatusResponse.class) .publishOn(asyncScheduler) .doOnSuccess(this::logBody) .map(this::doSomething); } } // 单元测试代码 public class StatusServiceUnitTest { private final RestClient mockClient = RestClient.builder() .exchangeFunction(req -> Mono.just( ClientResponse.create(HttpStatus.OK) .header("Content-Type", "application/json") .body("{\"status\":\"SUCCESS\"}") .build() )).build(); // 测试用immediate调度器,所有操作同步执行 private final StatusService service = new StatusService(mockClient, Schedulers.immediate(), "test-url"); @Test void shouldReturnProcessedResultWhenRequestSuccess() { StepVerifier.create(service.getProcessedStatus()) .expectNextMatches(result -> "PROCESSED_SUCCESS".equals(result.getStatus())) .verifyComplete(); } }
集成测试实现方案
集成测试需要验证完整流程,包括真实的调度器逻辑、网络调用逻辑:
- 用
WireMock启动本地模拟服务,代替真实的外部接口,返回预设的响应内容。 - 直接注入生产环境使用的真实
ASYNC_SCHEDULER,不需要替换,验证调度器的执行逻辑符合预期。 - 还是用
StepVerifier等待流执行完成,不需要手动加Thread.sleep,StepVerifier会自动订阅直到流终止,可自定义超时时间避免测试无限等待。
集成测试示例代码:
@WireMockTest(httpPort = 8090) public class StatusServiceIntegrationTest { private final RestClient restClient = RestClient.create(); // 注入生产环境真实调度器 private final StatusService service = new StatusService(restClient, AppConfig.ASYNC_SCHEDULER, "http://localhost:8090"); @Test void shouldRunProcessLogicOnAsyncScheduler() { // 配置WireMock返回响应 stubFor(get(urlEqualTo("/id")) .willReturn(okJson("{\"status\":\"SUCCESS\"}"))); StepVerifier.create(service.getProcessedStatus()) .expectNextMatches(result -> "PROCESSED_SUCCESS".equals(result.getStatus())) .verify(Duration.ofSeconds(2)); // 验证HTTP请求确实发送 verify(1, getRequestedFor(urlEqualTo("/id"))); } }
反应式多线程场景测试核心要点
- 永远不要用
Thread.sleep等待异步结果,StepVerifier的verify方法会自动等待流结束,内置超时配置可以避免阻塞。 - 调度器全部做成可注入的,单元测试用
Schedulers.immediate()或者虚拟时间调度器,既提升测试运行速度,也避免异步线程带来的测试不稳定问题。 - 如果需要验证逻辑确实运行在指定调度器上,可以在
doOnSuccess/doOnNext中记录当前线程名,测试时断言线程名符合调度器的命名前缀规则即可。 - 异常测试和同步场景写法一致,直接用
StepVerifier的expectError系列方法即可自动捕获多线程抛出的异常。
内容的提问来源于stack exchange,提问作者Melad Basilius
相关产品推荐
相关产品推荐

