Java如何实现异步调用队列,满足API限流条件时按序执行请求
修正后的实现可以满足你按顺序执行、限流排队的要求,核心是用单线程调度器全局消费请求队列,统一控制执行时机,对外暴露的接口和原始API完全兼容:
import java.util.concurrent.*; public class Sandbox2 { public static void main(String[] args) throws ExecutionException, InterruptedException { // 示例配置:两次请求最小间隔1000ms,即每秒最多1次请求,可按需调整 RateLimitedApiWrapper apiWrapper = new RateLimitedApiWrapper(1000); CompletableFuture<Integer> result1 = apiWrapper.requestAThing("Req1"); CompletableFuture<Integer> result2 = apiWrapper.requestAThing("Req2"); CompletableFuture<Integer> result3 = apiWrapper.requestAThing("Req3"); System.out.println("Result1: " + result1.get()); System.out.println("Result2: " + result2.get()); System.out.println("Result3: " + result3.get()); // 用完关闭调度器释放线程资源 apiWrapper.shutdown(); } // 限流API包装类 public static class RateLimitedApiWrapper { private final ActualApi actualApi = new ActualApi(); // 存储请求任务的阻塞队列 private final BlockingQueue<RequestTask> taskQueue = new LinkedBlockingQueue<>(); // 单线程调度器,保证请求严格按入队顺序执行 private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); // 两次请求的最小间隔,单位毫秒 private final long requestIntervalMs; public RateLimitedApiWrapper(long requestIntervalMs) { this.requestIntervalMs = requestIntervalMs; // 启动后台队列消费线程 scheduler.submit(this::processQueue); } // 对外暴露的请求方法,和原始API签名完全一致 public CompletableFuture<Integer> requestAThing(String reqParam) { CompletableFuture<Integer> resultFuture = new CompletableFuture<>(); // 请求直接入队,立即返回未完成的Future taskQueue.add(new RequestTask(reqParam, resultFuture)); return resultFuture; } // 后台队列消费逻辑 private void processQueue() { while (!scheduler.isShutdown()) { try { // 阻塞等待下一个请求 RequestTask task = taskQueue.take(); // 执行实际API调用 CompletableFuture<Integer> actualFuture = actualApi.requestHandler(task.reqParam); // 把实际调用的结果、异常转发到对外返回的Future actualFuture.whenComplete((res, ex) -> { if (ex != null) { task.resultFuture.completeExceptionally(ex); } else { task.resultFuture.complete(res); } }); // 限流等待,保证请求间隔符合阈值要求 Thread.sleep(requestIntervalMs); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } } public void shutdown() { scheduler.shutdown(); try { if (!scheduler.awaitTermination(60, TimeUnit.SECONDS)) { scheduler.shutdownNow(); } } catch (InterruptedException e) { scheduler.shutdownNow(); } } // 内部类,封装请求参数和对应的待完成Future private static class RequestTask { private final String reqParam; private final CompletableFuture<Integer> resultFuture; public RequestTask(String reqParam, CompletableFuture<Integer> resultFuture) { this.reqParam = reqParam; this.resultFuture = resultFuture; } } } // 原始API实现,和你的代码保持一致 public static class ActualApi { public CompletableFuture<Integer> requestHandler(String request) { return CompletableFuture.supplyAsync(() -> { try { Thread.sleep(1000); } catch (InterruptedException e) { e.printStackTrace(); } return Integer.parseInt(request.substring(3)); }); } } }
如果需要更灵活的限流策略(比如每分钟允许N次请求,允许小流量突发),可以用Guava的RateLimiter替换Thread.sleep逻辑即可,不需要修改其他代码。
内容的提问来源于stack exchange,提问作者mephisto_mori
相关产品推荐
相关产品推荐

