
在 Project Reactor 中,通过不可变数据结构传递状态并结合 onErrorMap 或 onErrorResume 捕获异常时的中间状态,可安全实现多阶段 API 调用中的状态维护与容错保存。
在 project reactor 中,通过不可变数据结构传递状态并结合 `onerrormap` 或 `onerrorresume` 捕获异常时的中间状态,可安全实现多阶段 api 调用中的状态维护与容错保存。
在响应式编程中,直接在 flatMap 链中复用并修改同一个可变对象(如示例中的 mDTO1)存在显著风险:不仅违背了函数式编程的无副作用原则,更在并发执行、重试(retry)、回压或取消等场景下极易引发状态竞争、数据覆盖或 NullPointerException。Reactors 的设计哲学强调不可变性(immutability)与纯函数式链式转换,因此推荐采用“每步生成新状态实例”的方式维护上下文。
✅ 推荐方案:使用不可变状态容器 + 异常增强
首先定义一个不可变的中间状态类(推荐使用 Lombok 的 @With 或记录类 record):
public record ProcessingState(
MyDTO dto,
Api1Response api1Response,
Api2Response api2Response,
Api3Response api3Response
) {
public static ProcessingState init(MyDTO dto) {
return new ProcessingState(dto, null, null, null);
}
public ProcessingState withApi1(Api1Response r) {
return new ProcessingState(dto, r, api2Response, api3Response);
}
public ProcessingState withApi2(Api2Response r) {
return new ProcessingState(dto, api1Response, r, api3Response);
}
public ProcessingState withApi3(Api3Response r) {
return new ProcessingState(dto, api1Response, api2Response, r);
}
}然后重构主流程,全程使用 ProcessingState 作为流转载体,并在每个 flatMap 中返回新实例:
public Mono<MyDTO> test(AnotherDTO req) {
MyDTO initialDto = new MyDTO(req);
ProcessingState initialState = ProcessingState.init(initialDto);
return Mono.just(initialState)
// Step 1: call api_1
.flatMap(state -> api_1.get(state.dto().getName())
.map(api1Res -> state.withApi1(api1Res))
.onErrorMap(e -> new StatefulException(state, e))) // ← 关键:包装当前状态
// Step 2: call api_2
.flatMap(state -> {
String var = state.api1Response().getSomeValue();
return api_2.get(var)
.map(api2Res -> state.withApi2(api2Res))
.onErrorMap(e -> new StatefulException(state, e));
})
// Step 3: call api_3
.flatMap(state -> {
String var1 = state.api1Response().getSomeValue2();
Integer var2 = state.api2Response().getSomeInt(); // 注意空安全
return api_3.get(var1, var2)
.map(api3Res -> state.withApi3(api3Res))
.onErrorMap(e -> new StatefulException(state, e));
})
// Final: persist full success
.flatMap(state -> Mono.fromRunnable(() -> saveIntoDB(state.dto())))
// Global error handler — 可访问失败前的完整状态
.onErrorResume(StatefulException.class, ex -> {
log.error("Failed at step with partial state: {}", ex.state(), ex.getCause());
// ✅ 此处可安全保存部分完成的状态
savePartialState(ex.state());
return Mono.empty();
})
.onErrorResume(Throwable.class, ex -> {
log.error("Unexpected error", ex);
return Mono.empty();
})
.map(state -> state.dto()); // 最终返回 MyDTO
}自定义异常类用于携带上下文状态:
public class StatefulException extends RuntimeException {
private final ProcessingState state;
public StatefulException(ProcessingState state, Throwable cause) {
super(cause);
this.state = state;
}
public ProcessingState state() {
return state;
}
}⚠️ 注意事项与最佳实践
-
禁止共享可变对象:原始代码中在多个
flatMap中反复修改mDTO1是反模式。Reactor 运算符可能被调度到不同线程,且flatMap内部默认允许并发(concurrency=256),导致竞态。 -
空值安全必须显式处理:
api1Response()等字段可能为null,调用.getSomeValue()前应判空或使用Optional封装。 -
onErrorMapvsonErrorResume:onErrorMap仅转换异常类型,不中断流;onErrorResume可替换为新流(如保存 + 继续)。二者结合使用更灵活。 -
替代简洁写法(适用于简单场景):若不想引入新类型,也可用
Tuple或Map<string object></string>临时承载,但语义性差、易出错,不推荐生产使用。 -
日志与监控建议:在
onErrorResume中记录state.toString()和异常堆栈,配合分布式追踪(如 Sleuth)可精准定位失败环节与上下文。
通过上述方式,你不仅能确保线程安全与响应式语义一致性,还能在任意环节失败时精确获取“截至上一步已完成的所有响应”,真正实现可观测、可恢复、可审计的响应式业务流。


















