Thread.interrupt对Java ArrayBlockingQueue::put的影响及队列清理问题
问题背景
使用ArrayBlockingQueue从Supplier向多个Worker传输工作项,通过Thread.interrupt通知Supplier停止提供新任务,但中断后清理队列的操作无法可靠清空队列。
清理队列函数
private static void clear(Thread supplier) { supplier.interrupt(); queue.clear(); }
Supplier核心代码
try { while (true) { Runnable r = () -> { /* 任务逻辑 */}; queue.put(r); } } catch (InterruptedException iex) { }
核心疑问
关于Thread.interrupt与ArrayBlockingQueue::put的交互存在疑问:
- 若队列已满且Supplier已携带中断标志,
put会等待还是抛出异常? - 若队列未满,Supplier在
put时被中断或已携带中断标志,会抛出异常还是修改队列?
对比ReentrantLock::lockInterruptibly的文档,其明确线程进入方法时已中断或获取锁时被中断,都会抛出InterruptedException并清除中断状态。而ArrayBlockingQueue内部使用ReentrantLock,put调用lockInterruptibly,clear调用lock。按此逻辑,中断后Supplier应无法获取锁修改队列,但测试代码反复出现断言失败,队列清理不可靠。
测试代码
import java.util.*; import java.util.concurrent.*; public class ABQ { private static void work() { while(true) { try { queue.take().run(); } catch (InterruptedException iex) { // 顶层中断 -> 不处理 } } } private static void clear(Thread caller, int level) { caller.interrupt(); queue.clear(); tripped.set(level); } private static ArrayBlockingQueue<Runnable> queue; private static ThreadLocal<Integer> tripped = ThreadLocal.withInitial(() -> -1); public static void main(String[] args) { queue = new ArrayBlockingQueue<>(5); Thread[] workers = new Thread[5]; for (int i=0; i < 5; ++i) { workers[i] = new Thread(ABQ::work); workers[i].start(); } final Thread caller = Thread.currentThread(); final Random rdm = new Random(); for (int i=0; ;i++) { final int finalI = i; try { Thread.sleep(2000); System.err.println("2s sleep over"); } catch (InterruptedException iex) { continue; } try { while (true) { Runnable r = () -> { int waitTime; boolean success; if (tripped.get() >= finalI) { throw new RuntimeException("Assertion failed, work for "+finalI+" but already tripped "+tripped.get()); } synchronized (rdm) { waitTime = rdm.nextInt(1000) + 200; success = rdm.nextBoolean(); } try { Thread.sleep(waitTime); if (success) { clear(caller, finalI); } } catch (InterruptedException iex) { // 内部中断 -> 同样忽略 } }; queue.put(r); } } catch (InterruptedException iex) { System.err.println("EOL"); } } } }
测试逻辑:反复填充队列直到被中断,进入下一批次。Worker会断言若发现已结束批次的工作项则报错,测试中所有Worker均触发断言失败。
原因分析
1. 中断与put操作的竞态问题
调用supplier.interrupt()后,Supplier线程可能处于以下阶段,导致清理不彻底:
- 已获取队列锁:如果Supplier已经通过
lockInterruptibly()拿到队列锁,此时即使被中断,put方法会先完成入队操作,之后才会响应中断抛出异常。这会导致queue.clear()执行前,已有新元素被加入队列。 - 中断与清理的时序差:
interrupt()和queue.clear()之间没有同步机制,Supplier可能在中断后、清理前完成一次put操作,导致队列残留旧批次任务。
2. ThreadLocal的误用
测试代码中tripped使用ThreadLocal,每个Worker线程持有独立副本。当某个Worker调用clear设置tripped时,仅能修改自身线程的副本,其他Worker无法感知批次结束的状态。这会导致部分Worker执行旧批次任务时无法触发断言,而当Supplier进入新批次后,残留的旧任务被执行时,若Worker自身的tripped已被设置,就会触发断言失败。
3. ArrayBlockingQueue.put的中断响应细节
根据ArrayBlockingQueue源码,put方法的中断响应逻辑如下:
public void put(E e) throws InterruptedException { checkNotNull(e); final ReentrantLock lock = this.lock; lock.lockInterruptibly(); try { while (count == items.length) notFull.await(); enqueue(e); } finally { lock.unlock(); } }
- 若线程调用
lockInterruptibly()前已被中断,会立即抛出InterruptedException,不执行入队。 - 若线程已获取锁,此时被中断:
- 队列未满时,会直接完成入队,解锁后退出,不会抛出异常(中断状态被保留),下一次循环调用
put时才会触发中断异常。 - 队列已满时,会在
notFull.await()时响应中断,抛出异常并清除中断状态。
- 队列未满时,会直接完成入队,解锁后退出,不会抛出异常(中断状态被保留),下一次循环调用
这意味着中断后,Supplier可能还会完成一次put操作,导致队列残留元素。
解决方案
1. 引入显式停止标志,避免竞态
使用AtomicBoolean作为停止标志,让Supplier主动检查停止信号,而非仅依赖中断:
private static AtomicBoolean stopSupplier = new AtomicBoolean(false); // 修改Supplier代码 try { while (!stopSupplier.get()) { Runnable r = () -> { /* 任务逻辑 */}; queue.put(r); } } catch (InterruptedException iex) { Thread.currentThread().interrupt(); // 保留中断状态,便于后续处理 } // 修改clear方法 private static void clear(Thread supplier) { stopSupplier.set(true); supplier.interrupt(); // 唤醒可能在put等待的线程 queue.clear(); stopSupplier.set(false); // 重置标志,供下一批次使用 }
停止标志确保Supplier主动退出循环,避免中断后仍完成入队操作。
2. 修复状态共享方式,替换ThreadLocal
将tripped改为全局原子变量,让所有Worker共享批次结束状态:
private static AtomicInteger tripped = new AtomicInteger(-1); // 修改Runnable中的断言逻辑 if (tripped.get() >= finalI) { throw new RuntimeException("Assertion failed, work for "+finalI+" but already tripped "+tripped.get()); } // 修改clear方法中的设置逻辑 tripped.set(level);
这样所有Worker都能实时感知批次结束的信号,准确触发断言。
3. 等待Supplier完全停止后再清理队列
若Supplier是独立线程,可在清理前等待其退出:
private static void clear(Thread supplier) { stopSupplier.set(true); supplier.interrupt(); try { supplier.join(); // 等待Supplier线程退出循环 } catch (InterruptedException iex) { Thread.currentThread().interrupt(); } queue.clear(); stopSupplier.set(false); }
注:此方法不适用于Supplier为主线程的场景,此时优先使用停止标志方案。
4. 使用drainTo处理残留任务
若需要处理被清理的任务(如取消),可使用drainTo替代clear:
List<Runnable> remainingTasks = new ArrayList<>(); queue.drainTo(remainingTasks); // 遍历remainingTasks,执行任务取消或其他清理逻辑
该方法原子性地清空队列,同时可获取残留任务进行后续处理。
总结
队列清理不可靠的核心是中断与put操作的竞态和ThreadLocal的状态共享错误。通过引入显式停止标志、修复状态共享机制、确保清理操作的原子性,可彻底解决该问题。
内容的提问来源于stack exchange,提问作者alex berne

