如何实现每秒最多5次请求的速率限制?含补时等待逻辑
接口限流实现方案(每秒最多5次请求,分组补时)
针对你的需求——对integrationClient.getOutageDetails()接口实现每秒最多5次请求,且当5次请求总耗时小于1秒时,等待差值时间再发起下一组请求,推荐基于批次处理+CountDownLatch的改造方案,贴合你现有代码结构,具体实现如下:
核心思路
- 将请求按5个为一组进行批次处理,确保同一时间最多并发5次请求;
- 记录每个批次的开始时间,用
CountDownLatch等待批次内所有请求完成; - 计算批次总耗时,若小于1000ms,则休眠剩余时间后再启动下一批次,严格控制每秒最多5次的频率。
代码改造示例
替换你原有for (Incident incidentOrChangeReq : listIncidentsOrChangeRequests)循环部分,修改后的完整代码片段如下:
int concurrentCalls = 5; ScheduledExecutorService executorService = Executors.newScheduledThreadPool(concurrentCalls); AtomicLong time = new AtomicLong(0); // 新增批次控制相关变量 int batchSize = 5; int totalRequests = listIncidentsOrChangeRequests.size(); int totalBatches = (totalRequests + batchSize - 1) / batchSize; // 向上取整计算总批次 try { for (int batchIdx = 0; batchIdx < totalBatches; batchIdx++) { long batchStartTime = System.currentTimeMillis(); // 初始化计数器,值为当前批次的实际请求数(最后一批可能不足5个) int currentBatchSize = Math.min(batchSize, totalRequests - batchIdx * batchSize); CountDownLatch batchLatch = new CountDownLatch(currentBatchSize); // 处理当前批次的请求 int startIdx = batchIdx * batchSize; int endIdx = Math.min(startIdx + batchSize, totalRequests); for (int i = startIdx; i < endIdx; i++) { Incident incidentOrChangeReq = listIncidentsOrChangeRequests.get(i); semaphore.acquire(); executorService.submit(() -> { Object listDetailsObject = new Object(); try { final StopWatch outageDetailsTime = new StopWatch(); outageDetailsTime.start(); // 目标限流接口 listDetailsObject = integrationClient.getOutageDetails(accountId, incidentOrChangeReq.getNumber()); outageDetailsTime.stop(); time.addAndGet(outageDetailsTime.getTotalTimeMillis()); log.error("outage details time: " + outageDetailsTime.getTotalTimeMillis()); final StopWatch setJsonTime = new StopWatch(); setJsonTime.start(); String detailsListJSON = ow.writeValueAsString(listDetailsObject); JSONObject jsonObject1 = new JSONObject(detailsListJSON); setValues(typeOfExport, jsonObject1, mapper, incidentOrChangeReq); setJsonTime.stop(); } catch (Exception e) { log.error(e.getMessage()); } finally { semaphore.release(); batchLatch.countDown(); // 任务完成后计数器减1 } }); } // 等待当前批次所有请求完成 try { batchLatch.await(); } catch (InterruptedException e) { log.error("等待批次任务中断: " + e.getMessage()); Thread.currentThread().interrupt(); break; } // 计算批次耗时,不足1秒则休眠补时 long batchCost = System.currentTimeMillis() - batchStartTime; if (batchCost < 1000) { long waitTime = 1000 - batchCost; try { Thread.sleep(waitTime); log.info("批次耗时不足1秒,等待{}ms后启动下一批次", waitTime); } catch (InterruptedException e) { log.error("补时等待中断: " + e.getMessage()); Thread.currentThread().interrupt(); break; } } } } catch (InterruptedException e) { log.error("InterruptedException caught: " + e.getMessage()); Thread.currentThread().interrupt(); } catch (Exception e) { log.error(e.getMessage()); }
关键细节说明
- Semaphore配置:确保
semaphore的许可数等于concurrentCalls(即5),保证同一时间最多5个并发请求; - CountDownLatch作用:精准等待批次内所有任务完成,避免提前计算批次耗时;
- 中断处理:所有中断场景都保留线程中断状态,避免后续逻辑异常;
- 最后批次兼容:自动处理最后一批请求不足5个的情况,无需额外判断。
备选方案(平滑限流)
如果不需要严格按5个一组,而是希望每秒最多5次的平滑限流,可以使用Guava的RateLimiter:
// 初始化限流器,每秒5个许可 RateLimiter rateLimiter = RateLimiter.create(5.0); // 在调用接口前获取许可 rateLimiter.acquire(); listDetailsObject = integrationClient.getOutageDetails(accountId, incidentOrChangeReq.getNumber());
但该方案无法保证“5次请求总耗时不足1秒则补时”的分组要求,仅适合平滑限流场景。
内容的提问来源于stack exchange,提问作者None
相关产品推荐
相关产品推荐

