如何在总并发14限制下实现多服务间动态线程分配?
问题描述
我有四个服务需向仅允许14个并发请求的远程服务器提交四类交易。其中Service A和Service B处理高量数据,各需至少6个线程;Service C和Service D各需至少1个线程。
规则要求
- 所有服务活跃时,每个服务需获得其最小线程配额;
- 若服务处于闲置状态,其分配的线程需可供其他服务使用以最大化资源利用率;
- 闲置服务激活时,需动态回收其分配的线程;
- 无论线程如何分配,总并发请求数绝不能超过14。
额外信息
- Service A和B每日运行一次,处理数据库中的大数据集并生成两类交易;
- Service C和D为按需触发服务,由用户请求触发并生成不同类型的交易。
我尝试使用以ServiceType为键、持有服务最小许可数的Semaphore为值的ConcurrentHashMap结合ExecutorService实现,代码如下:
public class TaskExecutorManagerImpl implements TaskExecutorManager { private final ExecutorService executorService; private final Map<ServiceType, Boolean> activeServices = new ConcurrentHashMap<>(); private final Map<ServiceType, Semaphore> serviceSemaphores = new ConcurrentHashMap<>(); public TaskExecutorManagerImpl() { executorService = new ThreadPoolExecutor( 0, ServiceType.getCount(), 30, TimeUnit.SECONDS, new SynchronousQueue<>(), new ThreadPoolExecutor.CallerRunsPolicy() ); for (ServiceType serviceType: ServiceType.values()) { serviceSemaphores.put(serviceType, new Semaphore(serviceType.getMinThreads())); activeServices.put(serviceType, false); } } @Override public void activateService(ServiceType serviceType) { activeServices.replace(serviceType, true); } @Override public void inactivateService(ServiceType serviceType) { activeServices.replace(serviceType, false); } @Override public void exec(ServiceType serviceType, Runnable task) { if (!activeServices.get(serviceType)) throw new RuntimeException(serviceType + " is inactive. Please activate it first."); Semaphore servicePermits = acquirePermit(serviceType); executorService.execute(() -> { try { task.run(); } finally { servicePermits.release(); } }); } @Override public <T> Future<T> submit(ServiceType serviceType, Callable<T> callable) { if (!activeServices.get(serviceType)) throw new RuntimeException(serviceType + " is inactive. Please activate it first."); Semaphore servicePermits = acquirePermit(serviceType); return executorService.submit(() -> { try { return callable.call(); } finally { servicePermits.release(); } }); } private Semaphore acquirePermit(ServiceType serviceType) { Semaphore servicePermits = serviceSemaphores.get(serviceType); while (true) { if (servicePermits.tryAcquire()) { return servicePermits; } else { Optional<ServiceType> inactiveService = activeServices.entrySet().stream() .filter(entry -> !entry.getValue()) .map(Map.Entry::getKey) .filter(s -> serviceSemaphores.get(s).availablePermits() > 0) .findFirst(); if (inactiveService.isPresent()) { Semaphore semaphore = serviceSemaphores.get(inactiveService.get()); if (semaphore.tryAcquire()) return semaphore; } } } } }
但测试发现该实现存在bug:会执行14个及以上任务后进入阻塞状态。我不一定必须使用ExecutorService,请问该如何正确实现需求?
解决方案
原代码核心问题分析
- 线程池配置错误:线程池最大线程数设为服务数量(4),远小于允许的14并发,结合
CallerRunsPolicy会导致调用线程直接执行任务,实际并发数不受控,甚至超过14。 - 无全局并发控制:仅靠各服务独立的Semaphore无法保证总并发不超过14,且闲置服务的许可借用逻辑没有和全局限制联动。
- 许可获取逻辑缺陷:
acquirePermit的无限循环在无可用许可时会空转,且没有处理服务激活时的配额回收逻辑。
正确实现方案
通过全局总并发Semaphore+服务配额跟踪的方式实现需求,同时调整线程池配置以匹配并发上限:
import java.util.Map; import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; public class TaskExecutorManagerImpl implements TaskExecutorManager { // 全局总并发限制:严格控制在14以内 private final Semaphore globalSemaphore = new Semaphore(14); // 跟踪每个服务已使用的许可数,确保活跃服务保留最小配额 private final Map<ServiceType, AtomicInteger> usedPermits = new ConcurrentHashMap<>(); // 线程池:最大线程数匹配总并发上限,避免调用线程执行任务导致并发失控 private final ExecutorService executorService = new ThreadPoolExecutor( 0, 14, 30, TimeUnit.SECONDS, new SynchronousQueue<>(), new ThreadPoolExecutor.AbortPolicy() // 超出并发时直接拒绝,可根据需求调整策略 ); // 服务活跃状态标记 private final Map<ServiceType, Boolean> activeServices = new ConcurrentHashMap<>(); public TaskExecutorManagerImpl() { for (ServiceType serviceType : ServiceType.values()) { activeServices.put(serviceType, false); usedPermits.put(serviceType, new AtomicInteger(0)); } } @Override public void activateService(ServiceType serviceType) { activeServices.replace(serviceType, true); int minThreads = serviceType.getMinThreads(); AtomicInteger serviceUsed = usedPermits.get(serviceType); // 激活时确保服务获得最小配额:先从其他活跃服务回收超额配额,再等待释放 while (serviceUsed.get() < minThreads) { boolean reclaimed = false; // 遍历其他活跃服务,回收其超出最小配额的许可 for (ServiceType other : ServiceType.values()) { if (other == serviceType || !activeServices.get(other)) { continue; } AtomicInteger otherUsed = usedPermits.get(other); if (otherUsed.get() > other.getMinThreads()) { otherUsed.decrementAndGet(); serviceUsed.incrementAndGet(); reclaimed = true; break; } } if (!reclaimed) { // 无法回收时,等待全局许可释放 try { globalSemaphore.acquire(); serviceUsed.incrementAndGet(); globalSemaphore.release(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("激活服务被中断", e); } } } } @Override public void inactivateService(ServiceType serviceType) { activeServices.replace(serviceType, false); // 闲置时释放所有已使用配额,供其他服务借用 usedPermits.get(serviceType).set(0); } @Override public void exec(ServiceType serviceType, Runnable task) { if (!activeServices.get(serviceType)) { throw new RuntimeException(serviceType + " 未激活,请先激活"); } try { acquirePermit(serviceType); executorService.execute(() -> { try { task.run(); } finally { releasePermit(serviceType); } }); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("执行任务被中断", e); } } @Override public <T> Future<T> submit(ServiceType serviceType, Callable<T> callable) { if (!activeServices.get(serviceType)) { throw new RuntimeException(serviceType + " 未激活,请先激活"); } try { acquirePermit(serviceType); return executorService.submit(() -> { try { return callable.call(); } finally { releasePermit(serviceType); } }); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("提交任务被中断", e); } } private void acquirePermit(ServiceType serviceType) throws InterruptedException { int minThreads = serviceType.getMinThreads(); AtomicInteger serviceUsed = usedPermits.get(serviceType); globalSemaphore.acquire(); try { // 未达最小配额时,直接占用全局许可 if (serviceUsed.get() < minThreads) { serviceUsed.incrementAndGet(); } else { // 已达最小配额,优先借用闲置服务的配额 boolean borrowed = false; for (ServiceType inactive : ServiceType.values()) { if (activeServices.get(inactive)) { continue; } AtomicInteger inactiveUsed = usedPermits.get(inactive); if (inactiveUsed.get() < inactive.getMinThreads()) { inactiveUsed.incrementAndGet(); borrowed = true; break; } } // 无闲置配额可借时,直接占用自身超额许可 if (!borrowed) { serviceUsed.incrementAndGet(); } } } catch (Exception e) { globalSemaphore.release(); throw e; } } private void releasePermit(ServiceType serviceType) { AtomicInteger serviceUsed = usedPermits.get(serviceType); // 优先归还借用的闲置服务配额 boolean returned = false; for (ServiceType inactive : ServiceType.values()) { if (activeServices.get(inactive)) { continue; } AtomicInteger inactiveUsed = usedPermits.get(inactive); if (inactiveUsed.get() > 0) { inactiveUsed.decrementAndGet(); returned = true; break; } } // 无借用配额可归还时,释放自身配额 if (!returned) { serviceUsed.decrementAndGet(); } globalSemaphore.release(); } } // 假设的ServiceType枚举定义 enum ServiceType { A(6), B(6), C(1), D(1); private final int minThreads; ServiceType(int minThreads) { this.minThreads = minThreads; } public int getMinThreads() { return minThreads; } public static int getCount() { return values().length; } } // 接口定义 interface TaskExecutorManager { void activateService(ServiceType serviceType); void inactivateService(ServiceType serviceType); void exec(ServiceType serviceType, Runnable task); <T> Future<T> submit(ServiceType serviceType, Callable<T> callable); }
实现说明
- 全局并发控制:
globalSemaphore(14)严格限制总并发数,确保永远不会超过14。 - 服务配额保障:通过
usedPermits跟踪每个服务的已使用许可,活跃服务始终保留其最小线程配额。 - 闲置配额复用:活跃服务需要额外线程时,优先借用闲置服务的配额;释放时优先归还借用的配额,最大化资源利用率。
- 激活回收逻辑:闲置服务激活时,自动从其他活跃服务的超额配额中回收资源,确保自身达到最小配额要求。
- 线程池适配:最大线程数设为14,匹配总并发上限,避免调用线程执行任务导致的并发失控。
内容的提问来源于stack exchange,提问作者Themisba7
相关产品推荐
相关产品推荐

