为何ThreadPoolExecutor代码无法终止?排查疑问与代码示例
我在IntelliJ IDEA中运行代码时,程序始终无法终止。改用exec.shutdownNow()则能成功停止,怀疑问题出在PrimeProducer(代码里的generator)上,但找不到具体原因。
主方法代码
public static void main(String[] args) { PrimeProducer generator = new PrimeProducer(new ArrayBlockingQueue<>(10)); ExecutorService exec = Executors.newFixedThreadPool(1); exec.execute(generator); try { Thread.sleep(1000); } catch (InterruptedException e) { throw new RuntimeException(e); } finally { generator.cancel(); } //generator.get(); exec.shutdown(); }
最初的PrimeProducer类
public class PrimeProducer extends Thread{ private final BlockingQueue<BigInteger> queue; PrimeProducer(BlockingQueue<BigInteger> queue){ this.queue = queue; } @Override public void run() { try { BigInteger p = BigInteger.ONE; while(!Thread.currentThread().isInterrupted()){ queue.put(p = p.nextProbablePrime()); } } catch (InterruptedException e) { } } public void cancel(){interrupt();} public void get(){ for(BigInteger i : queue){ System.out.println(i.toString()); } } }
仍无法正常运行的最新版本
public class PrimeProducer implements Runnable { private final BlockingQueue<BigInteger> queue; PrimeProducer(BlockingQueue<BigInteger> queue){ this.queue = queue; } @Override public void run() { try { BigInteger p = BigInteger.ONE; while(!Thread.currentThread().isInterrupted()){ queue.put(p = p.nextProbablePrime()); } } catch (InterruptedException e) { } } public void cancel(){Thread.currentThread().interrupt();} public synchronized void get(){ for(BigInteger i : queue){ System.out.println(i.toString()); } } }
问题根源
1. 最初版本的问题
当PrimeProducer继承Thread时,你通过exec.execute(generator)把它作为Runnable提交给线程池,此时generator对象本身的线程并没有被启动,真正执行run方法的是线程池中的线程。调用generator.cancel()时,实际中断的是generator这个未启动的Thread实例,而非线程池里正在运行任务的线程,中断信号根本没传到目标线程,导致while循环一直执行,程序无法终止。
2. 最新版本的问题
改成实现Runnable后,cancel()方法里调用Thread.currentThread().interrupt(),但这里的currentThread()是调用cancel方法的main线程,不是线程池里执行PrimeProducer任务的线程。你中断的是main线程,任务线程的循环依然不会停止,程序自然无法终止。
解决方案
方案一:用ExecutorService的原生中断机制
直接去掉自定义的cancel方法,在finally块中调用exec.shutdownNow(),它会向所有正在执行的任务线程发送中断信号,这也是你之前测试能成功停止的原因。如果想优雅等待一段时间再强制终止,可以结合awaitTermination:
finally { exec.shutdownNow(); // 可选:等待1秒看任务是否终止 try { if (!exec.awaitTermination(1, TimeUnit.SECONDS)) { // 若仍未终止,可添加额外处理逻辑 } } catch (InterruptedException e) { exec.shutdownNow(); } }
方案二:正确追踪任务线程并中断
如果一定要保留自定义cancel逻辑,需要在PrimeProducer中记录执行任务的线程:
public class PrimeProducer implements Runnable { private final BlockingQueue<BigInteger> queue; private volatile Thread workerThread; // 用volatile保证多线程下的可见性 PrimeProducer(BlockingQueue<BigInteger> queue){ this.queue = queue; } @Override public void run() { workerThread = Thread.currentThread(); // 记录当前执行任务的线程 try { BigInteger p = BigInteger.ONE; while(!Thread.currentThread().isInterrupted()){ queue.put(p = p.nextProbablePrime()); } } catch (InterruptedException e) { // 中断后直接退出循环即可 } finally { workerThread = null; // 清理线程引用 } } public void cancel(){ if (workerThread != null) { workerThread.interrupt(); // 中断真正执行任务的线程 } } public void get(){ for(BigInteger i : queue){ System.out.println(i.toString()); } } }
这样调用generator.cancel()时,就能正确中断线程池里执行任务的线程,while循环检测到中断状态后退出,任务结束后exec.shutdown()就能正常终止程序。
额外注意点
- 不要忽略
InterruptedException,如果后续还有逻辑需要处理,建议在catch块中恢复中断状态:Thread.currentThread().interrupt();,避免中断信号丢失。 - ArrayBlockingQueue是有界队列,当队列满时
queue.put()会阻塞,此时线程被中断会抛出InterruptedException,刚好能触发循环退出,这部分逻辑是正常的。
内容的提问来源于stack exchange,提问作者xiunian

