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

在Project Reactor响应式流中维护状态对象及错误时获取状态

使用Project Reactor维护响应式流中的状态对象

我正在开发一个调用多个API的响应式流序列,希望维护一个MyDTO状态对象。请问能否通过Project Reactor实现该需求?以下是我的实现代码及诉求:

public Mono<MyDTO> test(AnotherDTO req) {
        MyDTO mDTO1 = new MyDTO(req); // --> line_1
        // ... some code to update myDTO1 based on req object
        return api_1.get(mDTO1.getName()).flatMap(api_1_Res -> { // --> 1
            mDTO1.setApi1Response(api_1_Res);
            return Mono.just(mDTO1);
        }).flatMap(monoResp -> {
            // ... some code to get a value from monoResp
            String var = monoResp.getApi1Response().getSomeValue();
            return api_2.get(var).flatMap(api_2_Res -> {
                monoResp.setApi2Response(api_2_Res);
                return Mono.just(monoResp);
            });
        }).flatMap(monoResp -> {
            // ... some code to get a value from monoResp
            String var1 = monoResp.getApi1Response().getSomeValue2();
            String var2 = monoResp.getApi2Response().getSomeInt();
            return api_3.get(var1, var2).flatMap(api_3_Res -> {
                monoResp.setApi3Response(api_3_Res);
                return Mono.just(monoResp);
            });
        }).flatMap(monoResp -> {
            saveIntoDB(monoResp);
        }).onErrorResume(errRes -> {
            log.error("exception {}", errRes.getMessage());
            return Mono.empty();
        });
    }

问题

  • 如何维护状态对象,使得在任意flatMap操作(调用api_1、api_2、api_3时)发生错误时,MyDTO能保留此前操作的状态并在onErrorResume中访问?
  • 原问题补充:当前实现能否维护MyDTO状态?若api_2出错,如何获取包含api_1响应的MyDTO状态并存入数据库?

解决方案

问题分析

当前代码直接修改外部创建的MyDTO对象,这种外部可变状态在Reactor的异步非阻塞模型中存在线程安全风险——多个订阅者或并发线程可能同时修改该对象,导致状态不一致。此外,全局onErrorResume无法直接获取错误发生时的MyDTO当前状态,因为外部对象的修改时机不确定。

正确实现方式

核心思路是将状态对象流转在响应式流内部,避免外部可变状态,并在每个API调用步骤添加局部错误处理,确保错误发生时能捕获当前已更新的状态。

重构代码(可变MyDTO场景)

public Mono<MyDTO> test(AnotherDTO req) {
    // 将状态初始化放入流内部,避免外部可变
    return Mono.just(new MyDTO(req))
            // 根据req更新初始状态
            .map(mDTO -> {
                // ... 此处添加基于req更新mDTO的代码
                return mDTO;
            })
            // 调用api_1,更新状态并处理局部错误
            .flatMap(mDTO -> 
                api_1.get(mDTO.getName())
                    .map(api1Res -> {
                        mDTO.setApi1Response(api1Res);
                        return mDTO;
                    })
                    // api_1出错时,捕获当前仅含初始数据的mDTO状态
                    .onErrorResume(err -> {
                        log.error("api_1调用失败: {}", err.getMessage());
                        saveIntoDB(mDTO);
                        return Mono.empty();
                    })
            )
            // 调用api_2,基于api_1结果更新状态
            .flatMap(mDTO -> {
                String var = mDTO.getApi1Response().getSomeValue();
                return api_2.get(var)
                        .map(api2Res -> {
                            mDTO.setApi2Response(api2Res);
                            return mDTO;
                        })
                        // api_2出错时,捕获已包含api_1响应的mDTO状态
                        .onErrorResume(err -> {
                            log.error("api_2调用失败: {}", err.getMessage());
                            saveIntoDB(mDTO);
                            return Mono.empty();
                        });
            })
            // 调用api_3,基于api_1、api_2结果更新状态
            .flatMap(mDTO -> {
                String var1 = mDTO.getApi1Response().getSomeValue2();
                String var2 = mDTO.getApi2Response().getSomeInt();
                return api_3.get(var1, var2)
                        .map(api3Res -> {
                            mDTO.setApi3Response(api3Res);
                            return mDTO;
                        })
                        // api_3出错时,捕获已包含api_1、api_2响应的mDTO状态
                        .onErrorResume(err -> {
                            log.error("api_3调用失败: {}", err.getMessage());
                            saveIntoDB(mDTO);
                            return Mono.empty();
                        });
            })
            // 所有步骤成功后保存完整状态
            .doOnNext(this::saveIntoDB)
            // 全局兜底错误处理
            .onErrorResume(err -> {
                log.error("全局流错误: {}", err.getMessage());
                return Mono.empty();
            });
}

优化:使用不可变MyDTO(推荐)

如果将MyDTO设计为不可变类,每次状态更新都创建新实例,可彻底避免线程安全问题:

// 不可变MyDTO示例
public class MyDTO {
    private final AnotherDTO req;
    private final Api1Response api1Response;
    private final Api2Response api2Response;
    private final Api3Response api3Response;

    // 构造器
    public MyDTO(AnotherDTO req) {
        this(req, null, null, null);
    }

    public MyDTO(AnotherDTO req, Api1Response api1Response, Api2Response api2Response, Api3Response api3Response) {
        this.req = req;
        this.api1Response = api1Response;
        this.api2Response = api2Response;
        this.api3Response = api3Response;
    }

    // 状态更新方法,返回新实例
    public MyDTO withApi1Response(Api1Response api1Response) {
        return new MyDTO(this.req, api1Response, this.api2Response, this.api3Response);
    }

    public MyDTO withApi2Response(Api2Response api2Response) {
        return new MyDTO(this.req, this.api1Response, api2Response, this.api3Response);
    }

    public MyDTO withApi3Response(Api3Response api3Response) {
        return new MyDTO(this.req, this.api1Response, this.api2Response, api3Response);
    }

    // Getter方法...
}

对应的流操作修改为:

.map(api1Res -> mDTO.withApi1Response(api1Res))

关键说明

  1. 状态流转:所有状态的创建和更新都在响应式流内部完成,每个操作传递更新后的对象,避免外部可变状态。
  2. 局部错误捕获:每个API调用的onErrorResume能直接访问当前步骤的MyDTO状态,确保错误发生时拿到的是已完成的所有前置操作的结果(比如api_2出错时,状态已包含api_1的响应)。
  3. 线程安全:不可变MyDTO的方式从根本上消除了线程安全问题,是Reactor中维护状态的最佳实践。

内容的提问来源于Stack Exchange,提问作者Venu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 08:37:32