如何通过Mono与Flux限制HTTP请求的并发数量
问题原因
你使用window(15)没有生效的核心原因是:window操作符仅负责将流拆分为固定大小的窗口,没有控制内部请求并发的能力,而Flux.merge()默认会同时订阅所有传入的Mono,因此所有HTTP请求还是会同时发起。
解决方案
直接使用flatMap的内置并发参数即可实现固定并发数控制,flatMap的第二个参数concurrency可以指定最多同时订阅的内部Publisher数量,刚好符合你「最多同时15个待处理请求,完成一个自动补充一个」的需求。
修改后的代码示例
首先修改syncData方法,将Flux.merge(totalTask)替换为带并发限制的处理逻辑:
public Flux<JsonNode> syncData() { return service1 .getData(param1) .flatMapMany(res -> { List<Mono<JsonNode>> totalTask = new ArrayList<>(); Map<String, Object> originData = service2.getDataFromDB(param2); res.withArray("data").forEach(row -> { String id = row.get("id").asText(); if (originData.containsKey(id)) { totalTask.add(service1.updateRequest(param3)); } else { totalTask.add(service1.deleteRequest(param4)); } originData.remove(id); }); for (left : originData.keySet()) { // 补全你原有的遍历剩余数据的逻辑 totalTask.add(service1.createRequest(param5)); } // 替换原有的Flux.merge(totalTask) return Flux.fromIterable(totalTask) // 第二个参数指定最大并发数为15 .flatMap(task -> task, 15); }); }
如果你需要根据单个请求的返回结果触发新的请求,并且希望新请求也被纳入15并发的限制,可以使用expand操作符实现递归处理,示例逻辑如下:
return Flux.fromIterable(totalTask) .flatMap(task -> task, 15) // 递归处理返回结果生成的新请求,同样受15并发限制 .expand(res -> { List<Mono<JsonNode>> newTasks = generateNewTasksFromResponse(res); // 自行实现根据返回生成新请求的逻辑 return Flux.fromIterable(newTasks).flatMap(task -> task, 15); });
调用端示例
如果是最终消费流的方法,直接订阅即可:
void syncData() { syncDataService.syncData() // 处理每个返回结果 .doOnNext(res -> handleResponse(res)) // 处理异常 .doOnError(e -> log.error("同步失败", e)) // 触发订阅执行 .subscribe(); }
实现原理
flatMap的concurrency参数会严格限制同时活跃的内部订阅数:当同时运行的请求达到15个时,新的请求会进入排队队列,直到有正在运行的请求完成、内部订阅被释放,才会自动从队列中拉取新的请求执行,全程不需要手动控制窗口、也不会出现请求突增的情况。
内容的提问来源于stack exchange,提问作者zzangs33
相关产品推荐
相关产品推荐

