自定义Java线程池仅执行部分任务,请求排查与解决
问题排查与修复
核心问题
你遇到的情况是因为主线程提交任务后立即调用shutdown(),导致Worker线程在执行完第一个任务后就直接退出,没有处理队列中剩余的任务。具体原因有两点:
shutdown()会立刻给所有Worker线程设置isStopped = true,并中断等待的线程。- Worker的
run()方法中,执行完单个任务后会检查if(isStopped) return;,这使得线程在收到停止信号后,不会继续从队列中获取下一个任务,直接退出循环。
修复方案
我们需要调整线程池的停止逻辑,让Worker线程在收到停止信号后,先处理完队列中所有剩余任务,再退出线程。同时,主线程可以等待所有任务执行完毕再结束,避免程序提前退出。
修改后的完整代码
import java.util.LinkedList; import java.util.Queue; import java.util.concurrent.CountDownLatch; class Worker extends Thread { private final Queue<Runnable> taskQueue; private volatile boolean isStopped = false; private CountDownLatch latch; public Worker(Queue<Runnable> taskQueue, CountDownLatch latch) { this.taskQueue = taskQueue; this.latch = latch; } public void run() { // 循环条件改为:未停止 或 队列还有任务 while (!isStopped || !taskQueue.isEmpty()) { Runnable task = null; synchronized (taskQueue) { while (taskQueue.isEmpty() && !isStopped) { try { taskQueue.wait(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); // 保留中断标记,继续检查状态 } } if (!taskQueue.isEmpty()) { task = taskQueue.poll(); } } if (task != null) { try { task.run(); } finally { if (latch != null) { latch.countDown(); } } } } } public void stopThread() { isStopped = true; synchronized (taskQueue) { taskQueue.notifyAll(); // 唤醒等待线程,处理剩余任务 } } } public class CustomThreadPool { private final int poolSize; private final Queue<Runnable> taskQueue; private final Worker[] workers; private CountDownLatch taskLatch; public CustomThreadPool(int poolSize) { this.poolSize = poolSize; taskQueue = new LinkedList<>(); workers = new Worker[poolSize]; } public void init(int taskCount) { taskLatch = new CountDownLatch(taskCount); for (int i = 0; i < poolSize; i++) { workers[i] = new Worker(taskQueue, taskLatch); workers[i].start(); } } public void execute(Runnable task) { synchronized (taskQueue) { taskQueue.offer(task); taskQueue.notifyAll(); } } public void shutdown() { for (Worker worker : workers) { worker.stopThread(); } } public void awaitTermination() throws InterruptedException { if (taskLatch != null) { taskLatch.await(); } for (Worker worker : workers) { worker.join(); } } public static void main(String[] args) throws InterruptedException { int taskCount = 5; CustomThreadPool threadPool = new CustomThreadPool(3); threadPool.init(taskCount); for (int i = 1; i <= taskCount; i++) { int num = i; threadPool.execute(() -> { int result = calculateFactorial(num); System.out.println("Factorial of " + num + " is " + result + " - Thread: " + Thread.currentThread().getName()); }); } threadPool.shutdown(); threadPool.awaitTermination(); } private static int calculateFactorial(int n) { int factorial = 1; for (int i = 1; i <= n; i++) { factorial *= i; } return factorial; } }
关键修改点
- Worker线程循环条件:将
while (!isStopped)改为while (!isStopped || !taskQueue.isEmpty()),确保即使收到停止信号,线程也会处理完队列中所有剩余任务。 - 停止逻辑优化:
stopThread()中调用taskQueue.notifyAll(),唤醒所有等待的线程,让它们立刻检查队列中的剩余任务。 - 主线程等待机制:新增
CountDownLatch和awaitTermination()方法,让主线程等待所有任务执行完毕再退出,避免程序提前终止导致输出不完整。 - 中断处理优化:Worker线程被中断后,保留中断标记,继续检查队列状态,而不是直接退出。
运行结果
修改后运行代码,会输出所有5个任务的结果:
Factorial of 1 is 1 - Thread: Thread-0 Factorial of 2 is 2 - Thread: Thread-1 Factorial of 3 is 6 - Thread: Thread-2 Factorial of 4 is 24 - Thread: Thread-0 Factorial of 5 is 120 - Thread: Thread-1
内容的提问来源于stack exchange,提问作者Aditya Tiwari
相关产品推荐
相关产品推荐

