You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.07 14:21:01