
在 Project Reactor 中,直接复用可变状态对象(如 MyDTO)存在线程安全与副作用风险;推荐采用不可变状态传递或异常封装策略,在错误发生时仍能访问已构建的中间状态。
在 project reactor 中,直接复用可变状态对象(如 `mydto`)存在线程安全与副作用风险;推荐采用不可变状态传递或异常封装策略,在错误发生时仍能访问已构建的中间状态。
在响应式编程中,状态管理的核心原则是避免共享可变状态。你当前代码中在 flatMap 链中反复修改同一个 MyDTO 实例(如 mDTO1.setApi1Response(...)),虽在单线程、顺序执行的简单场景下看似可行,但一旦涉及异步切换、重试(retry)、并行组合(zip, flatMap 并发调用)或背压重放,就极易引发竞态条件、数据覆盖或 NullPointerException —— 因为 Mono/Flux 的操作符不保证执行线程,且 flatMap 内部可能被调度到不同线程池。
✅ 推荐方案一:不可变状态建模(首选)
将 MyDTO 改为不可变类(或使用 Builder 模式),每次更新都返回新实例。这天然线程安全,语义清晰,也契合函数式响应式编程范式:
// 不可变 DTO 示例(Lombok 简化)
@Value
public class MyDTO {
String name;
Api1Response api1Response;
Api2Response api2Response;
Api3Response api3Response;
// 提供 withXxx 方法(手动或 Lombok @With)
public MyDTO withApi1Response(Api1Response r) {
return new MyDTO(this.name, r, this.api2Response, this.api3Response);
}
}
// 响应式链重构后
public Mono<MyDTO> test(AnotherDTO req) {
MyDTO initial = new MyDTO(req.getName(), null, null, null);
return api_1.get(req.getName())
.map(Api1Response::new)
.map(initial::withApi1Response)
.flatMap(dto -> api_2.get(dto.getApi1Response().getSomeValue())
.map(Api2Response::new)
.map(dto::withApi2Response))
.flatMap(dto -> api_3.get(
dto.getApi1Response().getSomeValue2(),
dto.getApi2Response().getSomeInt())
.map(Api3Response::new)
.map(dto::withApi3Response))
.doOnNext(this::saveIntoDB) // 成功时落库
.onErrorResume(err -> {
log.error("Flow failed at step with partial state: {}", err.getMessage());
// 此处 err 发生时,dto 已丢失(因 flatMap 中断)→ 需结合 onErrorMap 捕获中间态
return Mono.empty();
});
}⚠️ 注意:上述链中,若 api_2.get(...) 抛异常,上游 dto(含 api1Response)无法直接访问 —— 因为 flatMap 未完成,dto 是局部变量。此时需配合 错误增强策略。
✅ 推荐方案二:错误时携带上下文(关键实践)
使用 onErrorMap 将原始异常包装为包含当前状态的自定义异常,确保错误传播链中始终携带“最后有效状态”:
public class StatefulException extends RuntimeException {
private final MyDTO partialState;
public StatefulException(MyDTO partialState, Throwable cause) {
super(cause);
this.partialState = partialState;
}
public MyDTO getPartialState() { return partialState; }
}
// 在每一步 flatMap 后添加状态快照与错误封装
public Mono<MyDTO> testWithStatefulError(AnotherDTO req) {
MyDTO initial = new MyDTO(req.getName(), null, null, null);
return Mono.just(initial)
.flatMap(dto -> api_1.get(dto.getName())
.map(Api1Response::new)
.map(dto::withApi1Response)
.onErrorMap(err -> new StatefulException(dto, err))) // ← 记录 api1 前状态
.flatMap(dto -> api_2.get(dto.getApi1Response().getSomeValue())
.map(Api2Response::new)
.map(dto::withApi2Response)
.onErrorMap(err -> new StatefulException(dto, err))) // ← 记录含 api1 的状态
.flatMap(dto -> api_3.get(
dto.getApi1Response().getSomeValue2(),
dto.getApi2Response().getSomeInt())
.map(Api3Response::new)
.map(dto::withApi3Response)
.onErrorMap(err -> new StatefulException(dto, err)))
.doOnNext(this::saveIntoDB)
.onErrorResume(StatefulException.class, ex -> {
MyDTO failedState = ex.getPartialState();
log.warn("Flow failed with partial state: {}", failedState);
savePartialStateToDB(failedState); // 可选:落库失败快照
return Mono.error(ex.getCause()); // 或返回 Mono.empty()
})
.onErrorResume(Throwable.class, err -> {
log.error("Unexpected error", err);
return Mono.empty();
});
}? 额外建议:
- 避免在
flatMap内部直接调用阻塞方法(如saveIntoDB若为同步 JDBC),应使用Mono.fromCallable(...).subscribeOn(Schedulers.boundedElastic())。 - 对于复杂编排,可考虑
StepVerifier单元测试验证各阶段状态与错误路径。 - 若需跨多个
Mono共享状态,可使用Context(但仅限传递只读元数据,不适用于大对象或频繁更新)。
总之,不可变性 + 错误封装 是 Reactor 中安全维护业务状态的黄金组合:既保障并发正确性,又赋予错误处理完整的上下文感知能力。


















