从Flux特定元素创建Mono:当前实现是否为不良实践?
这种做法属于典型的不良实践,存在多个严重问题:
核心问题分析
- Mono永远无法正常完成:你的代码依赖
doOnCancel触发monoSink.success(),但takeUntil在条件满足时是正常终止订阅(触发onComplete),而非取消订阅。这意味着当结果集收集完成后,monoSink根本不会调用success(),下游会一直等待结果,永远收不到响应。 - 线程安全隐患:使用非线程安全的
HashMap存储结果,如果Flux是多线程推送元素,putAll操作会引发并发修改异常,或者导致数据不一致。 - 资源泄漏:直接调用
subscribe()却未保存订阅的Disposable,也没将内部订阅与Mono的生命周期绑定。当下游取消Mono订阅时,内部的Flux订阅不会被销毁,会持续消耗资源处理共享Flux的后续元素。 - 逻辑不直观、维护成本高:把Mono的完成逻辑和Flux的取消事件绑定,不符合Reactor声明式编程的设计思路,代码逻辑绕弯,后期调试和修改容易出错。
更合理的实现方式
用Reactor原生操作符实现,避免手动管理MonoSink和订阅:
private <T> Mono<Map<String, Optional<T>>> createTargetMono(Flux<Map<String, Optional<T>>> flux, List<String> queryParams) { // 使用线程安全的ConcurrentHashMap避免并发问题 return flux .reduce(new ConcurrentHashMap<>(), (result, currentMap) -> { // 累积需要的子集数据 result.putAll(subset(currentMap, queryParams)); return result; }) // 当结果集包含所有目标参数时终止 .takeUntil(result -> result.keySet().containsAll(queryParams)) // 转换为Mono(takeUntil会发出满足条件的第一个结果后终止) .next(); }
这个实现利用reduce累积数据,takeUntil控制终止时机,完全遵循Reactor的编程模型,避免了手动管理订阅和Sink的风险,同时保证了线程安全和资源自动回收。
内容的提问来源于stack exchange,提问作者Siri
相关产品推荐
相关产品推荐

