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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 01:27:04