线程池Worker执行顺序及忙碌状态任务分配问题求助
解决线程池Worker输出顺序与忙碌状态标记问题
看起来你在实现线程池时遇到了两个核心问题:输出顺序混乱,以及Worker忙碌状态的正确标记与任务分配。我来帮你一步步解决这些问题:
一、解决输出顺序混乱:用同步屏障确保启动信息全部输出后再执行任务
你遇到的输出顺序问题,本质是线程调度的不确定性——每个Worker线程启动后会立即开始执行任务,系统调度可能让某个Worker先执行任务,而其他Worker还没输出启动信息。要让所有Worker先输出"已启动"的信息,再统一开始执行任务,我们可以用CountDownLatch来做同步屏障:
- 在创建Worker的线程池类中,初始化一个
CountDownLatch,计数等于Worker的总数。 - 每个Worker启动后,先输出启动信息,然后调用
countDown()减少计数器,再调用await()等待所有Worker都完成启动信息的输出,之后再进入任务执行逻辑。
这样就能保证所有启动信息都打印完毕后,才会开始执行任务,输出顺序就不会混乱了。
二、正确标记Worker的忙碌状态,实现空闲Worker分配任务
原代码里构造函数将busy初始化为true是错误的——刚创建的Worker还没有任务,应该处于空闲状态(busy=false)。我们需要调整忙碌状态的时机:
- 初始化:Worker构造函数中设置
busy = false。 - 分配任务时:遍历Worker列表,找到第一个
busy=false的Worker,为它设置任务,并将busy设为true(注意要加锁,避免多线程下的竞态条件)。 - 任务执行完成后:Worker执行完任务后,将
busy设为false,同时释放资源、记录任务完成信息,等待下一次分配任务。
三、修改后的Worker.java代码示例
import java.util.ArrayList; import java.util.Map; import java.util.TreeMap; import java.util.concurrent.CountDownLatch; public class Worker implements Runnable { private final int id; private final JobStack jobStack; private final ResourceStack resourceStack; private Job job; private Resource[] resources; private boolean busy; // 空闲时为false,忙碌时为true private final Map<Integer, ArrayList<Integer>> jobsCompleted; private final CountDownLatch startLatch; // 同步启动的屏障 // 构造函数新增CountDownLatch参数 public Worker(int theId, JobStack theJobStack, ResourceStack theResourceStack, CountDownLatch startLatch) { id = theId; jobStack = theJobStack; resourceStack = theResourceStack; job = null; busy = false; // 初始化为空闲状态 jobsCompleted = new TreeMap<>(); this.startLatch = startLatch; } @Override public void run() { try { // 第一步:输出启动信息 System.out.printf("Worker %d 已启动工作%n", id); // 通知线程池:本Worker已完成启动信息输出 startLatch.countDown(); // 等待所有Worker都输出启动信息 startLatch.await(); // 第二步:进入任务循环,处理任务 while (true) { // 加锁避免多线程分配任务时的竞态条件 synchronized (this) { // 如果当前没有任务且处于空闲,等待任务分配 while (!busy && job == null) { wait(); } } if (job == null) { // 没有任务了,退出循环 break; } try { // 输出开始执行任务的信息 System.out.printf("Worker %d 开始执行任务 %d%n", id, job.getId()); // 模拟任务执行(替换成你的实际任务逻辑) Thread.sleep(1000); // 记录任务完成信息 ArrayList<Integer> resourceIds = new ArrayList<>(); for (Resource r : resources) { resourceIds.add(r.getId()); } jobsCompleted.put(job.getId(), resourceIds); // 任务完成后释放资源 resourceStack.releaseResources(resources); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } finally { // 任务完成后,标记为空闲,清空任务和资源 synchronized (this) { job = null; resources = null; busy = false; notifyAll(); // 通知线程池本Worker已空闲 } } } } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.printf("Worker %d 被中断%n", id); } } // 分配任务的方法,由线程池调用 public synchronized boolean assignJob(Job newJob, Resource[] newResources) { if (!busy) { this.job = newJob; this.resources = newResources; this.busy = true; notify(); // 唤醒Worker线程开始执行任务 return true; } return false; // 当前Worker忙碌,分配失败 } // 按需添加getter方法 public int getId() { return id; } public boolean isBusy() { synchronized (this) { return busy; } } }
四、线程池类的关键逻辑示例
import java.util.ArrayList; import java.util.List; import java.util.concurrent.CountDownLatch; public class ThreadPool { private final List<Worker> workers; private final JobStack jobStack; private final ResourceStack resourceStack; public ThreadPool(int workerCount, JobStack jobStack, ResourceStack resourceStack) { this.jobStack = jobStack; this.resourceStack = resourceStack; this.workers = new ArrayList<>(); // 创建启动同步的CountDownLatch CountDownLatch startLatch = new CountDownLatch(workerCount); // 创建并启动所有Worker for (int i = 0; i < workerCount; i++) { Worker worker = new Worker(i + 1, jobStack, resourceStack, startLatch); workers.add(worker); new Thread(worker).start(); } // 启动任务分配线程 new Thread(this::assignJobs).start(); } private void assignJobs() { while (true) { Job job = jobStack.getNextJob(); if (job == null) { // 所有任务分配完毕,通知所有Worker退出 for (Worker worker : workers) { synchronized (worker) { worker.notify(); } } break; } // 获取任务所需资源 Resource[] resources = resourceStack.getResources(job.getResourceCount()); if (resources == null) { // 资源不足,等待资源释放 try { Thread.sleep(500); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } continue; } // 找到第一个空闲的Worker分配任务 boolean assigned = false; while (!assigned) { for (Worker worker : workers) { if (worker.assignJob(job, resources)) { assigned = true; break; } } // 如果所有Worker都忙碌,等待一段时间再重试 if (!assigned) { try { Thread.sleep(300); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } } } } }
关键说明
- 同步屏障:
CountDownLatch确保了所有Worker先输出启动信息,再统一进入任务执行阶段,彻底解决输出顺序混乱的问题。 - 忙碌状态管理:所有对
busy的访问都加了synchronized锁,避免多线程下的竞态条件;Worker空闲时会调用wait()等待任务,分配任务时用notify()唤醒。 - 任务分配逻辑:线程池的
assignJobs方法会遍历Worker列表,只给空闲的Worker分配任务,确保忙碌的Worker不会被重复分配任务。
内容的提问来源于stack exchange,提问作者melih anik
相关产品推荐
相关产品推荐

