如何正确链式调用Mono/Flux 移除响应式代码中的block()阻塞调用
响应式链式调用正确实现方案
你之前的实现存在两个核心问题:
flatMap内部调用stateService.getStatesByIds()后没有将返回的Mono作为lambda的返回值,导致该流未被纳入整体订阅链路,不会被触发执行- 空值逻辑没有和流的终止信号绑定,直接在外层返回
Mono.just(mutation)会无视State查询结果,固定返回Mutation对象
重构后无阻塞的getMutation实现
完全对齐原业务逻辑,移除所有阻塞调用:
private Mono<Mutation> getMutation(int gameId, int typeId, int domainId) { return mutationProvider.findByGameIdAndTypeId(gameId, typeId) // Mutation不存在时直接抛错,和原逻辑一致 .switchIfEmpty(Mono.error(() -> new RuntimeException("Mutation was not found by gameId and typeId"))) // 串联State查询逻辑,必须将内部Mono作为flatMap返回值才能被订阅执行 .flatMap(mutation -> stateService.getStatesByIds(mutation.getId(), domainId) // 过滤State为空、或State.value为空的场景,不满足条件时直接返回空Mono .filter(state -> state != null && state.getValue() != null) // 空值场景打印日志,和原逻辑一致 .switchIfEmpty(Mono.fromRunnable(() -> log.info("Requested mutation with gameId[{}] typeId[{}] domainId[{}] is disabled. Value is null.", gameId, typeId, domainId) )) // 校验通过时组装Mutation属性,最终返回有效对象 .flatMap(state -> { mutation.setTemplateId(state.getTemplateId()); return Mono.just(mutation); }) ); }
配套调整说明
- 移除Service层不必要的调度器配置
重构后所有逻辑都是纯响应式非阻塞调用,没有需要隔离的阻塞代码,原来的.subscribeOn(Schedulers.boundedElastic())可以直接删除,避免无意义的线程池开销,调整后Service层代码:public Flux<Mutation> getMutations(List<MutationRequest> mutationRequests, int domainId) { return Flux.fromIterable(mutationRequests) .flatMap(mutationRequest -> getMutation(mutationRequest.getGameId(), mutationRequest.getTypeId(), domainId)); } - 关键注意事项
- 响应式流中所有中间操作生成的
Mono/Flux必须作为操作符lambda的返回值层层传递,最终被框架层的订阅者拿到才会触发执行,没有被返回的响应式对象属于悬空流,不会执行任何逻辑 filter操作符会在条件不满足时直接返回空Mono,自动终止当前链路的元素传递,不需要手动写if判断返回空值;上层Flux.flatMap会自动忽略空Mono,不会向下游SSE流推送无效数据- 非边缘场景禁止在响应式链路内部调用
block()/blockFirst()方法,这类方法会强行阻塞当前线程等待结果,完全违背响应式编程的非阻塞原则
- 响应式流中所有中间操作生成的
内容的提问来源于stack exchange,提问作者ayemaman
相关产品推荐
相关产品推荐

