如何在不终止resSink时,通过reduce将Flux转为Mono并及时返回结果
问题原因
你的代码中getRes方法使用reduce操作,而reduce必须等待上游Flux(即resSink.asFlux())完成才会发射最终累加结果。但resSink是一个永久存活的Sink(没有调用emitComplete()),所以reduce会一直阻塞,导致返回的Mono永远不会完成,测试也就永久挂起。
解决方案
替换reduce为scan+filter+take(1)组合:
scan会在每次收到新数据时发射当前的累加结果,而非等Flux结束才发射最终结果filter检查累加结果是否已经包含当前请求所需的所有keytake(1)拿到第一个符合条件的结果后立即终止流,让Mono正常完成
同时修复startListen中的嵌套订阅问题(Reactor反模式),改用flatMap处理异步逻辑。
修改后的服务代码
@Slf4j public class TestSinkService { public static final int MAX_BUFFER_SIZE = 2; public static final Duration MAX_BUFFER_TIME = Duration.ofMillis(500); private Sinks.Many<List<String>> reqSink = Sinks.many().replay().all(); private Sinks.Many<Map<String, Double>> resSink = Sinks.many().replay().all(); public void startListen() { reqSink .asFlux() .bufferTimeout(MAX_BUFFER_SIZE, MAX_BUFFER_TIME) .flatMap(lists -> { List<String> collect = lists.stream() .flatMap(List::stream) .distinct() .collect(Collectors.toList()); log.info("Fetching from client {}", collect); return getDataFromClientDumb(collect); }) .subscribe(stringDoubleMap -> { log.info("Received from client {}", stringDoubleMap); resSink.emitNext(stringDoubleMap, Sinks.EmitFailureHandler.FAIL_FAST); }, error -> { log.error("Failed to fetch data", error); }); } public Mono<Map<String, Double>> fetchData(List<String> list) { Mono<Map<String, Double>> mapMono = getRes(list); reqSink.emitNext(list, Sinks.EmitFailureHandler.FAIL_FAST); return mapMono; } private Mono<Map<String, Double>> getDataFromClientDumb(List<String> list) { Map<String, Double> data = list.stream().collect(Collectors.toMap(Function.identity(), v -> Math.random())); return Mono.just(data); } private Mono<Map<String, Double>> getRes(List<String> list) { Set<String> requiredKeys = new HashSet<>(list); return resSink.asFlux() .scan(new ConcurrentHashMap<>(), (accu, next) -> { log.info("Next {}", next); for (String req : list) { if (next.containsKey(req)) { accu.put(req, next.get(req)); } } log.info("Accu {}", accu); return accu; }) .filter(accu -> accu.keySet().containsAll(requiredKeys)) .take(1) .single(); } }
关键修改说明
getRes方法的核心优化:- 用
HashSet存储需要的key,方便快速判断是否收集完成 scan替代reduce,每次收到新结果就更新累加器并发射当前状态filter判断累加器是否包含所有请求的key,满足条件就放行take(1)拿到第一个符合条件的结果后立即终止流,确保Mono正常完成- 使用
ConcurrentHashMap避免多线程场景下的并发修改问题
- 用
startListen的嵌套订阅修复:- 把原来嵌套在
subscribe里的getDataFromClientDumb调用移到flatMap中,保持Reactor链式调用的风格 - 新增错误处理逻辑,避免异常导致整个流终止
- 把原来嵌套在
测试代码说明
修改后的服务代码会在收集到当前请求所需的所有数据后立即返回结果,StepVerifier能正常收到onComplete信号,不会再出现永久挂起的情况。
内容的提问来源于stack exchange,提问作者santik
相关产品推荐
相关产品推荐

