Java生产者-消费者模型实现中的死锁问题求助
生产者-消费者实现中的死锁问题
我尝试用2个消费者线程(Consumer1打印偶数、Consumer2打印奇数)和1个生产者线程实现生产者-消费者模型,虽然知道用了反模式且存在错误,但仅作为实验用途。目前每次运行程序都疑似死锁,控制台无任何输出。
类信息
- NumbersThread:继承体系顶层类,提供共享Queue
- ProducerThread:随机生成0到49范围内的数字并推入Queue
- Consumer1Thread:打印Queue中的偶数(Consumer2逻辑一致,负责打印奇数)
代码实现
NumbersThread
class NumbersThread extends Thread { private int capacity = 64; private final Queue<Integer> numbersQueue; // field variables, constructors etc. // peek, poll, and add methods synchronized on numbersQueue protected final void notifyAllForFull() { synchronized (FULL_LOCK) { FULL_LOCK.notifyAll(); } } protected final void notifyAllForEmpty() { synchronized (EMPTY_LOCK) { EMPTY_LOCK.notifyAll(); } } protected final void waitOnEmpty() throws InterruptedException { synchronized (EMPTY_LOCK) { EMPTY_LOCK.wait(); } } protected final void waitOnFull() throws InterruptedException { synchronized (FULL_LOCK) { FULL_LOCK.wait(); } }
ProducerThread
class ProducerThread extends NumbersThread { ProducerThread(Queue<Integer> numbersQueue) { super(numbersQueue); } @Override public void run() { run: while (true) { while (isFull()) { try { waitOnFull(); } catch (InterruptedException e) { break run; } } var rand = new Random(); add(rand.nextInt(50)); notifyAllForEmpty(); } }
Consumer1Thread(Consumer2逻辑一致)
class Consumer1Thread extends NumbersThread { Consumer1Thread(Queue<Integer> numbersQueue) { super(numbersQueue); } @Override public void run() { run: while (true) { while (isEmpty()) { try { waitOnEmpty(); } catch (InterruptedException e) { break run; } } Integer num = peek(); if (num != null && num % 2 == 0 ) { int n = poll(); notifyAllForFull(); System.out.println("Consumer 1 (even): " + n); } } }
问题分析
- 无意义的空循环阻塞:当队列中的元素不符合当前消费者的处理类型时(比如Consumer1遇到奇数),线程不会释放锁也不会进入等待状态,而是持续循环检查队列头部元素,导致队列锁被长期占用,生产者和另一个消费者无法执行操作,最终所有线程陷入卡死。
- 等待通知机制失效:当前的wait/notify依赖EMPTY_LOCK和FULL_LOCK,但消费者在无法处理元素时没有正确触发等待,破坏了线程间的协作逻辑。
- 线程安全间隙:peek操作后到poll操作之间存在时间窗口,可能出现元素被其他线程修改的情况,同时空循环会持续消耗CPU资源。
修复方案
1. 修正消费者逻辑,添加条件等待
当消费者遇到无法处理的元素时,释放队列锁并等待新元素到来,避免持续占用锁。修改后的Consumer1Thread示例:
class Consumer1Thread extends NumbersThread { Consumer1Thread(Queue<Integer> numbersQueue) { super(numbersQueue); } @Override public void run() { run: while (true) { synchronized (numbersQueue) { while (isEmpty()) { try { waitOnEmpty(); } catch (InterruptedException e) { break run; } } Integer num = peek(); // 循环等待直到找到符合条件的元素或队列变空 while (num != null && num % 2 != 0) { try { numbersQueue.wait(); // 释放锁,等待新元素通知 } catch (InterruptedException e) { break run; } num = isEmpty() ? null : peek(); if (num == null) break; } if (num != null) { int n = poll(); notifyAllForFull(); System.out.println("Consumer 1 (even): " + n); numbersQueue.notifyAll(); // 通知其他消费者检查新元素 } } } } }
2. 调整生产者的通知逻辑
生产者添加元素后,除了通知等待空队列的线程,还需要通知所有消费者线程,让它们检查新元素是否符合自身处理条件:
@Override public void run() { run: while (true) { synchronized (numbersQueue) { while (isFull()) { try { waitOnFull(); } catch (InterruptedException e) { break run; } } var rand = new Random(); add(rand.nextInt(50)); notifyAllForEmpty(); numbersQueue.notifyAll(); // 触发所有消费者检查新元素 } } }
3. 统一锁粒度
确保所有队列相关操作(add、peek、poll、isEmpty、isFull)都在同一个锁(numbersQueue)下执行,避免因锁不一致导致的线程安全问题。
内容的提问来源于stack exchange,提问作者Snefru Clone
相关产品推荐
相关产品推荐

