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

多线程批量执行异常排查:为何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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 02:22:07