多线程批量执行异常排查:为何10线程未按3批并发运行?
问题:多线程批量执行逻辑不符合预期的原因分析
提供的代码
Client类
public class Client { Integer id; Integer priority; // 补充原代码缺失的构造方法,否则CPUDemo无法正常实例化Client public Client(Integer id, Integer priority) { this.id = id; this.priority = priority; } }
CPU类
package target2024.systemDesign.cpuProcessor; import lombok.SneakyThrows; import java.util.LinkedList; import java.util.Queue; //Singleton design pattern public class CPU { private static CPU instance; private final Integer MAX_THREAD_POOL = 3; Object lock = new Object(); static Queue<Client> executingThreadPool; static Queue<Client> waitingThreadPool; private CPU() { executingThreadPool = new LinkedList<>(); waitingThreadPool = new LinkedList<>(); } public synchronized static CPU getInstance() { if(instance == null) { instance = new CPU(); } return instance; } @SneakyThrows public void execute(Client client) { System.out.println("Received client=" + client.id); waitingThreadPool.add(client); synchronized (lock) { if(executingThreadPool.size() < MAX_THREAD_POOL) { executingThreadPool.add(waitingThreadPool.poll()); } else { lock.wait(); executingThreadPool.add(waitingThreadPool.poll()); } } Client clientToProcess; synchronized (lock) { clientToProcess = executingThreadPool.poll(); } process(clientToProcess); synchronized (lock) { lock.notifyAll(); } } @SneakyThrows public void process(Client client) { System.out.println("-----Executing client=" + client.id); Thread.sleep(1000); } }
CPUDemo类
package target2024.systemDesign.cpuProcessor; import lombok.SneakyThrows; public class CPUDemo { @SneakyThrows public static void main(String[] args) { CPU cpu = CPU.getInstance(); int clientSize = 10; Thread[] tarr = new Thread[clientSize]; Client[] carr = new Client[clientSize]; //Create clients for(int i=0; i<clientSize; i++) { carr[i] = new Client(i, i); } //Create threads for(int i=0; i<clientSize; i++) { int finalI = i; tarr[i] = new Thread(new Runnable() { @Override public void run() { cpu.execute(carr[finalI]); } }); } //Initialize threads for(int i=0; i<clientSize; i++) { tarr[i].start(); } } }
代码存在的核心问题
1. 执行逻辑完全偏离目标,未限制并发
execute方法的流程完全没有实现“同时最多3个任务执行”的限制:
- 每个线程进来后,把自己的
Client加入等待队列,紧接着就从等待队列取出放到执行队列,然后立刻从执行队列取出这个Client调用process。 - 这相当于每个线程都直接处理自己的任务,
executingThreadPool队列根本没起到限制并发的作用,所有线程的process会几乎同时启动。
2. wait/notify逻辑混乱,无法控制批次
- 当执行队列“满”时线程进入
wait,但被唤醒后直接取任务执行,没有检查此时是否真的有空位(多个线程被同时唤醒时,会导致执行任务数超过3个)。 - 任务完成后调用
notifyAll,但唤醒的线程会直接执行任务,无法实现“一批3个完成后再执行下一批”的批次控制,只是无限制的并发。
3. 执行队列设计完全失效
executingThreadPool的设计意图是跟踪正在执行的任务,但实际代码中,线程刚把任务放进去就立刻取出来执行,队列的size永远不会超过1,完全无法用来判断当前并发数。
4. 静态队列的冗余风险
executingThreadPool和waitingThreadPool被定义为static,虽然是单例模式,但静态成员属于类级别,会导致即使CPU实例被回收(单例场景下不会),队列仍占用资源,且增加了线程安全的维护成本。
修正思路(实现批次执行)
要实现“每次同时运行3个,完成后再执行下一批”,可以采用以下两种方式:
方式1:使用Semaphore+CountDownLatch实现批次控制
// 修改后的CPU类核心逻辑 public class CPU { private static CPU instance; private final int BATCH_SIZE = 3; private final Semaphore semaphore = new Semaphore(BATCH_SIZE); private CPU() {} public synchronized static CPU getInstance() { if(instance == null) { instance = new CPU(); } return instance; } @SneakyThrows public void executeBatch(List<Client> clients) { // 分批次处理任务 for (int i = 0; i < clients.size(); i += BATCH_SIZE) { int end = Math.min(i + BATCH_SIZE, clients.size()); List<Client> batch = clients.subList(i, end); CountDownLatch latch = new CountDownLatch(batch.size()); for (Client client : batch) { semaphore.acquire(); // 获取执行名额 new Thread(() -> { try { process(client); } finally { semaphore.release(); // 释放执行名额 latch.countDown(); // 标记当前任务完成 } }).start(); } latch.await(); // 等待当前批次所有任务完成 System.out.println("=== 当前批次完成,启动下一批 ==="); } } @SneakyThrows public void process(Client client) { System.out.println("-----Executing client=" + client.id); Thread.sleep(1000); } }
方式2:修正原wait/notify逻辑(基于原代码改造)
如果要保留原代码结构,需要重新设计队列和执行逻辑:
- 用计数器代替
executingThreadPool,跟踪当前正在执行的任务数量。 - 线程必须先获取执行名额才能调用
process,执行完成后释放名额并唤醒等待线程。 - 增加批次计数器,确保当前批次所有任务完成后再启动下一批。
内容的提问来源于stack exchange,提问作者Bhavuk Mathur
相关产品推荐
相关产品推荐

