基于多PriorityBlockingQueues的ThreadPoolExecutor:Java 8 Web请求调度执行
Java 8 分组式Web请求调度与配额管控实现方案
针对你在Java 8中处理大量分组式Web请求并管控配额的需求,我整理了一套实用的实现方案,结合Java 8的异步特性和线程安全工具来满足你的核心要求:分组不可变确定、配额跟踪与限制、高效请求调度执行。
核心思路拆解
要解决这个问题,我们需要三个核心组件:
- 请求模型:封装不可变分组标识与实际Web请求逻辑
- 配额管理器:线程安全地跟踪每个分组的配额使用情况
- 请求调度器:基于配额管控异步执行请求,利用Java 8的
CompletableFuture实现非阻塞处理
1. 定义不可变请求模型
首先我们需要一个WebRequest类,确保分组标识(比如用户ID)是不可变且确定的,同时封装实际的Web请求任务:
import java.util.concurrent.Callable; public class WebRequest { // 不可变的分组标识,比如用户ID、业务线ID等 private final String groupId; // 封装实际Web请求逻辑,比如调用HTTP接口、处理响应 private final Callable<String> requestTask; public WebRequest(String groupId, Callable<String> requestTask) { this.groupId = groupId; this.requestTask = requestTask; } // 仅提供只读访问,保证分组不可修改 public String getGroupId() { return groupId; } public Callable<String> getRequestTask() { return requestTask; } }
2. 线程安全的配额管理器
根据你的需求,配额管控分为两种常见场景:并发数限制(同一分组同时执行的请求数)和计数型配额(比如每日请求上限)。下面给出两种实现思路:
场景1:并发数配额管控(同一分组同时最多N个请求)
使用Semaphore控制每个分组的并发数,结合ConcurrentHashMap实现线程安全的分组存储:
import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.Semaphore; public class ConcurrentQuotaManager { private final int defaultConcurrentPermits; private final Map<String, Semaphore> groupSemaphores = new ConcurrentHashMap<>(); public ConcurrentQuotaManager(int defaultConcurrentPermits) { this.defaultConcurrentPermits = defaultConcurrentPermits; } // 懒加载创建分组对应的信号量 private Semaphore getGroupSemaphore(String groupId) { return groupSemaphores.computeIfAbsent(groupId, k -> new Semaphore(defaultConcurrentPermits)); } // 尝试获取配额,返回是否成功 public boolean tryAcquireQuota(String groupId) { return getGroupSemaphore(groupId).tryAcquire(); } // 执行完成后释放配额(无论请求成功/失败都要调用) public void releaseQuota(String groupId) { getGroupSemaphore(groupId).release(); } }
场景2:计数型配额管控(比如每日请求上限)
使用AtomicInteger跟踪每个分组的剩余配额,适合一次性或时间窗口内的限额:
import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicInteger; public class CountQuotaManager { private final int defaultDailyLimit; private final Map<String, AtomicInteger> groupQuotas = new ConcurrentHashMap<>(); public CountQuotaManager(int defaultDailyLimit) { this.defaultDailyLimit = defaultDailyLimit; } // 尝试获取配额,成功则扣减剩余次数 public boolean tryAcquireQuota(String groupId) { AtomicInteger quota = groupQuotas.computeIfAbsent(groupId, k -> new AtomicInteger(defaultDailyLimit)); // 仅当剩余配额>0时才扣减 return quota.updateAndGet(current -> current > 0 ? current - 1 : current) > 0; } // 可选:定期重置配额(比如每日零点),可配合ScheduledExecutorService实现 public void resetQuota(String groupId) { groupQuotas.put(groupId, new AtomicInteger(defaultDailyLimit)); } }
3. 分组请求调度器
结合Java 8的CompletableFuture和线程池,实现异步请求调度,同时集成配额管控逻辑:
import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; public class GroupedRequestScheduler { private final ConcurrentQuotaManager quotaManager; private final ExecutorService executor; public GroupedRequestScheduler(int defaultConcurrentQuota, int threadPoolSize) { this.quotaManager = new ConcurrentQuotaManager(defaultConcurrentQuota); // 根据实际请求量调整线程池类型,比如用newCachedThreadPool或自定义ThreadPoolExecutor this.executor = Executors.newFixedThreadPool(threadPoolSize); } // 提交请求并返回异步结果 public CompletableFuture<String> submitRequest(WebRequest request) { String groupId = request.getGroupId(); // 先校验配额,不足则直接返回失败的Future if (!quotaManager.tryAcquireQuota(groupId)) { return CompletableFuture.failedFuture(new RuntimeException("Quota exceeded for group: " + groupId)); } // 异步执行请求,最终释放配额 return CompletableFuture.supplyAsync(() -> { try { // 执行实际Web请求逻辑 return request.getRequestTask().call(); } catch (Exception e) { throw new RuntimeException("Request failed for group: " + groupId, e); } finally { // 无论请求成功还是失败,都要释放配额(并发配额场景) quotaManager.releaseQuota(groupId); } }, executor); } // 关闭调度器,释放线程池资源 public void shutdown() { executor.shutdown(); } }
4. 使用示例
下面是一个简单的测试示例,模拟同一分组的多个请求,验证配额管控效果:
public class RequestSchedulerDemo { public static void main(String[] args) throws InterruptedException { // 初始化调度器:每个分组最多同时执行2个请求,线程池大小10 GroupedRequestScheduler scheduler = new GroupedRequestScheduler(2, 10); // 创建用户A的3个测试请求(模拟耗时1秒的Web请求) WebRequest userAReq1 = new WebRequest("user_001", () -> { Thread.sleep(1000); return "[User 001] Request 1 completed"; }); WebRequest userAReq2 = new WebRequest("user_001", () -> { Thread.sleep(1000); return "[User 001] Request 2 completed"; }); WebRequest userAReq3 = new WebRequest("user_001", () -> { Thread.sleep(1000); return "[User 001] Request 3 completed"; }); // 提交请求并处理结果 scheduler.submitRequest(userAReq1).thenAccept(System.out::println); scheduler.submitRequest(userAReq2).thenAccept(System.out::println); scheduler.submitRequest(userAReq3) .thenAccept(System.out::println) .exceptionally(ex -> { System.err.println("Request failed: " + ex.getMessage()); return null; }); // 等待请求执行完成 Thread.sleep(3500); // 关闭调度器 scheduler.shutdown(); } }
关键注意事项
- 分组不可变性:确保
WebRequest的groupId是final且仅提供只读访问,避免分组被篡改 - 线程安全:所有分组相关的存储和操作都使用线程安全的容器(
ConcurrentHashMap)和原子类/信号量,避免多线程下的竞态问题 - 资源回收:请求执行完成后一定要释放配额(并发场景),避免资源泄漏;线程池使用完成后要调用
shutdown()释放资源 - 扩展能力:如果需要不同分组配置不同配额,可以修改配额管理器,支持传入分组-配额的初始配置;如果需要时间窗口配额(比如每秒N次),可以结合
ScheduledExecutorService定期重置配额
内容的提问来源于stack exchange,提问作者Ferenc Dósa-Rácz
相关产品推荐
相关产品推荐

