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

如何在总并发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,请问该如何正确实现需求?


解决方案

原代码核心问题分析

  1. 线程池配置错误:线程池最大线程数设为服务数量(4),远小于允许的14并发,结合CallerRunsPolicy会导致调用线程直接执行任务,实际并发数不受控,甚至超过14。
  2. 无全局并发控制:仅靠各服务独立的Semaphore无法保证总并发不超过14,且闲置服务的许可借用逻辑没有和全局限制联动。
  3. 许可获取逻辑缺陷: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);
}

实现说明

  1. 全局并发控制:globalSemaphore(14)严格限制总并发数,确保永远不会超过14。
  2. 服务配额保障:通过usedPermits跟踪每个服务的已使用许可,活跃服务始终保留其最小线程配额。
  3. 闲置配额复用:活跃服务需要额外线程时,优先借用闲置服务的配额;释放时优先归还借用的配额,最大化资源利用率。
  4. 激活回收逻辑:闲置服务激活时,自动从其他活跃服务的超额配额中回收资源,确保自身达到最小配额要求。
  5. 线程池适配:最大线程数设为14,匹配总并发上限,避免调用线程执行任务导致的并发失控。

内容的提问来源于stack exchange,提问作者Themisba7

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 01:05:55