Flux结合WebClient使用flatMap时如何实现固定批次并行请求速率控制
问题原因
你之前的buffer改造方案只是将元素按6个分组后立刻又拍平为连续流,下游flatMap的并发控制逻辑还是会在单个请求完成后立刻拉取新元素补充并行度,自然无法实现按批次等待的效果。
解决方案
调整算子顺序,将批次处理逻辑放到concatMap中实现,利用concatMap严格等待前一个Publisher完成再处理下一个的特性,实现批次间串行等待、批次内并行执行的效果,代码如下:
return Flux.fromIterable(new Generator()) .log() .buffer(6) // 按6个元素为一组拆分批次 // concatMap保证前一批所有请求处理完成后,才会拉取下一个批次 .concatMap(batch -> Flux.fromIterable(batch) .flatMap(s -> webClient .head() .uri( MessageFormat.format( "/my-{2,number,#00}.xml", channel, timestamp, s)) .exchangeToMono(r -> Mono.just(r.statusCode())) .filter(HttpStatus::is2xxSuccessful) .map(r -> s), 6 // 单批次内并行度为6,整批同时发起请求 ) ) .take(6) // 全局累计收集到6个成功响应后立刻终止流 .sort();
逻辑说明
- 上游
buffer(6)将源源不断的生成元素拆分为一个个大小为6的列表 concatMap收到一个批次列表后,会将其转为内部Flux,并行执行6个HEAD请求,等待该批次所有请求处理完成后,才会向上游请求下一个批次- 全局
take(6)会累计收集所有批次的成功响应,凑够6个后直接取消上游流,不再处理后续请求
内容的提问来源于stack exchange,提问作者Bojan Vukasovic
相关产品推荐
相关产品推荐

