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

Spring WebFlux:取两个Mono首个返回值,异步处理另一个并记录错误

Spring WebFlux:返回首个成功结果后保留慢流执行并记录错误

场景与需求

现有两个数据源,均返回Mono<Entity>类型的创建结果:

class CacheCustomerClient {
    Mono<Entity> createCustomer(Customer customer)
}

class MasterCustomerClient {
    Mono<Entity> createCustomer(Customer customer)
}

Spring WebFlux控制器接收创建请求:

@PostMapping
@ResponseStatus(HttpStatus.CREATED)
public Flux<Entity> createCustomer(@RequestBody Customer customer) {
    return customerService.createNewCustomer(customer);
}

核心需求:

  • 任一数据源创建成功,立即向调用者返回响应
  • 未完成的数据源需继续执行,若执行失败则记录错误日志

当前问题

使用Flux.firstWithValue()可以获取首个成功结果,但会向未完成的流传播取消信号,导致慢流被终止,错误无法被记录。

解决方案

核心思路是让两个创建流独立订阅执行,不受firstWithValue的取消逻辑影响。具体实现如下:

public Flux<Entity> createCustomer(final Customer customer) {
    // 定义缓存创建流,包含错误日志处理
    Mono<Entity> cacheCreate = cacheClient.createCustomer(customer)
            .doOnError(WebClientResponseException.class, err -> 
                log.error("缓存创建客户失败: {}", err.getMessage()));
    
    // 定义主库创建流,包含错误日志处理
    Mono<Entity> masterCreate = masterClient.createCustomer(customer)
            .doOnError(WebClientResponseException.class, err -> 
                log.error("主库创建客户失败: {}", err.getMessage()));

    // 独立订阅两个流,触发异步执行,不受后续取消信号影响
    cacheCreate.subscribe();
    masterCreate.subscribe();

    // 返回首个成功的结果,同时不中断另一个流的执行
    return Flux.firstWithValue(cacheCreate, masterCreate)
            .onErrorMap(err -> new RuntimeException("缓存和主库创建客户均失败", err));
}

关键说明

  1. 独立订阅的作用:调用subscribe()会单独触发流的执行,该订阅不受Flux.firstWithValue()的取消信号干扰,即使首个结果已返回,慢流仍会继续执行,错误日志能正常输出。
  2. 异常处理扩展:如果需要处理非WebClientResponseException的错误,可以给subscribe()添加全局错误处理器:
    cacheCreate.subscribe(
        result -> log.debug("缓存创建成功"),
        err -> log.error("缓存创建异常", err)
    );
    
  3. 调度器优化:若创建操作涉及阻塞逻辑,可通过subscribeOn()指定单独的调度器,避免阻塞主线程:
    Mono<Entity> cacheCreate = cacheClient.createCustomer(customer)
            .subscribeOn(Schedulers.boundedElastic())
            .doOnError(WebClientResponseException.class, err -> 
                log.error("缓存创建客户失败: {}", err.getMessage()));
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 02:06:04