Java生产者消费者项目报current thread is not owner异常排查
停车场生产者消费者模型线程唤醒异常排查
问题背景
开发大学课程项目,实现停车场管理模型:
- 生产者(Producer)线程:负责将车辆停入停车场
- 消费者(Consumer)线程:负责将车辆移出停车场
要求所有生产者、消费者运行在独立线程上,必须使用Mutex/Semaphores、Lock或Monitor机制实现停车场资源的互斥访问。
当前程序无法实现生产者与消费者的跨线程通知唤醒,运行时抛出如下异常:
java.lang.IllegalMonitorStateException: current thread is not owner
预期业务逻辑
- 停车场为空时,消费者调用
wait()进入休眠 - 若消费者取车前停车场已满,取车后唤醒所有休眠的生产者
- 若生产者停车前停车场为空,停车后唤醒所有休眠的消费者
项目原始源码
Consumer类
package prog; import java.util.Scanner; public class Consumer implements Runnable { private static int PARKING_GARAGE_SIZE = 4; Mutex consumerMutex; Mutex producerMutex; Buffer<Car> carBuffer; Thread currentThread; public Consumer(Buffer<Car> carBuffer, Mutex consumerMutex, Mutex producerMutex) { this.carBuffer = carBuffer; this.consumerMutex = consumerMutex; this.producerMutex = producerMutex; } public void run() { System.out.println("Consumer läuft"); int counter = 1; synchronized(this) { try { while (true) { System.out.println("Consumer-Schleife beginnt"); currentThread = Thread.currentThread(); String threadName = currentThread.getName() + counter++; consumerMutex.mutex.acquire(); System.out.println("Mutex acquired"); consumerMutex.WorkingQueue.add(currentThread); System.out.println("Consumers consumes on thread: " + threadName); if(carBuffer.full()) { System.out.println("Buffer full. Wake producers. Parke aus."); // wake producers producerMutex.WorkingQueue.notifyAll(); System.out.println("Producer geweckt"); carBuffer.pop(); consumerMutex.mutex.release(); consumerMutex.semaphore.release(); Thread.sleep((long)(Math.random() * 1500) + 500); } else if(carBuffer.empty()){ System.out.println("Buffer empty. Go to sleep."); wait(); consumerMutex.mutex.release(); consumerMutex.semaphore.release(); } else { System.out.println("Parke aus"); carBuffer.pop(); consumerMutex.mutex.release(); consumerMutex.semaphore.release(); Thread.sleep((long)(Math.random() * 1500) + 500); } System.out.println("Consumer-Schleife zu Ende"); } } catch (Exception e) { System.out.println("Consumer Catch-Error " + e.toString()); } } } public static void main(String[] args) { Buffer<Car> carBuffer = new Buffer<Car>(PARKING_GARAGE_SIZE); Mutex producerMutex = new Mutex(); Mutex consumerMutex = new Mutex(); Producer producer = new Producer(carBuffer, consumerMutex, producerMutex); Thread producerThread = new Thread(producer); Consumer consumer = new Consumer(carBuffer, consumerMutex, producerMutex); Thread consumerThread = new Thread(consumer); producerThread.start(); consumerThread.start(); } }
Producer类
package prog; import java.lang.Math; public class Producer implements Runnable { Mutex producerMutex; Mutex consumerMutex; Buffer<Car> carBuffer; Thread currentThread; public Producer(Buffer<Car> carBuffer, Mutex consumerMutex, Mutex producerMutex) { this.carBuffer = carBuffer; this.consumerMutex = consumerMutex; this.producerMutex = producerMutex; } public Mutex getProducerMutex() { return producerMutex; } public void run() { System.out.println("Producer läuft"); int counter = 1; synchronized(this) { try { while (true) { System.out.println("Producer-Schleife beginnt"); currentThread = Thread.currentThread(); String threadName = currentThread.getName() + counter++; producerMutex.mutex.acquire(); System.out.println("Mutex acquired"); producerMutex.WorkingQueue.add(currentThread); System.out.println("Producer produces on thread: " + threadName); Car newCar = new Car(); if(carBuffer.empty()) { System.out.println("Parkhaus leer. Wecke Consumer. Parke Auto Nummer " + (counter - 1)); consumerMutex.WorkingQueue.notifyAll(); notifyAll(); System.out.println("Consumer geweckt"); carBuffer.push(newCar); } else if(carBuffer.full()) { System.out.println("Buffer full, go sleep!"); wait(); } else{ System.out.println("Parke Auto Nummer " + (counter - 1)); carBuffer.push(newCar); } producerMutex.mutex.release(); producerMutex.semaphore.release(); Thread.sleep((long)(Math.random() * 1500) + 500); System.out.println("Producer-Schleife zu Ende"); } } catch (Exception e) { System.out.println("Producer Catch-Error: " + e.toString()); } } } }
Mutex类
package prog; import java.util.concurrent.Semaphore; import java.util.LinkedList; public class Mutex { LinkedList<Thread> WorkingQueue = new LinkedList<Thread>(); Semaphore semaphore = new Semaphore(0); Semaphore mutex = new Semaphore(1); }
Buffer类
package prog; import java.util.List; import java.util.Scanner; import java.util.ArrayList; public class Buffer<T> { private final ArrayList<T> bufferStack; private int size; private int top; private final Object FULL_QUEUE = new Object(); private final Object EMPTY_QUEUE = new Object(); public Buffer(int maxSize) { bufferStack = new ArrayList<>(maxSize); this.size = maxSize; top = -1; } public int getSize() { return size; } public boolean full() { return top == getSize() -1; } public void waitForFull() throws InterruptedException { } public boolean empty() { return top == -1; } public void push(T item) { if(!full()) { bufferStack.add(0, item); top++; } else { System.out.println("Error: Stack is full!"); } } public T pop() { if (!empty()) { T item = bufferStack.get(0); bufferStack.remove(0); top--; return item; } else { System.out.println("Error: Stack is empty!"); return null; } } public T peek() { if (!bufferStack.isEmpty()) { return bufferStack.get(0); } else { return null; } } public List<T> getCompleteStack() { return bufferStack; } }
Car类
package prog; import java.lang.Math; public class Car { private int Platenumber; public Car() { Platenumber = (int)(Math.random() * 999) + 1; } public int getPlatenumber() { return Platenumber; } }
问题根因
IllegalMonitorStateException 触发的唯一原因是:调用wait()/notify()/notifyAll()的线程,没有持有方法所属对象的监视器锁。当前代码的并发逻辑存在多处核心错误:
- 锁对象不匹配
run()方法中synchronized(this)持有的是Producer/Consumer实例自身的锁,但代码中调用producerMutex.WorkingQueue.notifyAll()、consumerMutex.WorkingQueue.notifyAll()时,从未获取过WorkingQueue对象的锁,直接调用notify方法必然抛出异常。此外无参wait()是在Producer/Consumer实例上等待,生产者和消费者持有的是完全独立的锁对象,根本不可能实现跨线程唤醒。 - 锁释放逻辑错误
缓冲区为空/满的分支中,代码先调用wait()让线程进入休眠,之后才写了mutex和semaphore的release逻辑。线程进入等待状态后不会执行后续代码,锁会被永久占用,其他线程永远无法获取锁访问临界区。 - 缺少循环状态校验
等待逻辑没有放在循环中,线程被唤醒后不会重新校验缓冲区状态,会出现虚假唤醒问题:比如消费者被唤醒时缓冲区仍为空,直接执行pop操作会引发数据错误。 - 临界区未被完全保护
缓冲区空/满判断、push/pop操作没有被同一个锁统一保护,多线程并发下会出现竞态条件,比如两个生产者同时判断缓冲区不满,写入时触发数组越界。
修正方案
不需要额外维护多余的Mutex类,直接以Buffer实例作为共享锁,将所有缓冲区操作、等待唤醒逻辑封装在Buffer内部,从根源避免锁不匹配问题。修正后的核心代码如下:
修正后的Buffer类
package prog; import java.util.ArrayList; import java.util.List; public class Buffer<T> { private final ArrayList<T> bufferStack; private final int size; private int top; public Buffer(int maxSize) { bufferStack = new ArrayList<>(maxSize); this.size = maxSize; top = -1; } public boolean full() { return top == size -1; } public boolean empty() { return top == -1; } // 同步停车方法 public synchronized void push(T item) throws InterruptedException { // 循环判断缓冲区是否满,满则生产者等待 while (full()) { wait(); } bufferStack.add(0, item); top++; // 停车后唤醒所有等待的消费者 notifyAll(); } // 同步取车方法 public synchronized T pop() throws InterruptedException { // 循环判断缓冲区是否空,空则消费者等待 while (empty()) { wait(); } T item = bufferStack.remove(0); top--; // 取车后唤醒所有等待的生产者 notifyAll(); return item; } public List<T> getCompleteStack() { return new ArrayList<>(bufferStack); } }
修正后的Producer类
package prog; import java.lang.Math; public class Producer implements Runnable { private final Buffer<Car> carBuffer; public Producer(Buffer<Car> carBuffer) { this.carBuffer = carBuffer; } public void run() { int counter = 1; try { while (true) { Car newCar = new Car(); System.out.println("Producer attempts to park car No." + counter); carBuffer.push(newCar); System.out.println("Successfully parked car No." + counter + ", current parking count: " + (carBuffer.getCompleteStack().size())); counter++; Thread.sleep((long)(Math.random() * 1500) + 500); } } catch (Exception e) { System.out.println("Producer error: " + e.toString()); } } }
修正后的Consumer类
package prog; public class Consumer implements Runnable { private static final int PARKING_GARAGE_SIZE = 4; private final Buffer<Car> carBuffer; public Consumer(Buffer<Car> carBuffer) { this.carBuffer = carBuffer; } public void run() { int counter = 1; try { while (true) { System.out.println("Consumer attempts to pick up car"); Car car = carBuffer.pop(); System.out.println("Successfully picked up car with plate number: " + car.getPlatenumber() + ", current parking count: " + carBuffer.getCompleteStack().size()); counter++; Thread.sleep((long)(Math.random() * 1500) + 500); } } catch (Exception e) { System.out.println("Consumer error: " + e.toString()); } } public static void main(String[] args) { Buffer<Car> carBuffer = new Buffer<>(PARKING_GARAGE_SIZE); Thread producerThread = new Thread(new Producer(carBuffer)); Thread consumerThread = new Thread(new Consumer(carBuffer)); producerThread.start(); consumerThread.start(); } }
逻辑说明
- 所有对缓冲区的操作都用
synchronized修饰,锁对象为Buffer实例本身,生产者消费者操作临界区时必须先获取这把锁,从根源避免锁不匹配问题 - 等待逻辑放在
while循环中,被唤醒后会重新检查缓冲区状态,避免虚假唤醒 - 每次push/pop操作完成后调用
notifyAll()唤醒所有等待的线程,由线程自己判断是否满足执行条件,完全符合预期业务逻辑 - 如果课程要求必须使用Semaphore实现,只需要将Buffer中的内置锁替换为1个互斥信号量、2个计数信号量(分别统计空车位、已占车位数量)即可,核心逻辑不变。
内容的提问来源于stack exchange,提问作者Clytax
相关产品推荐
相关产品推荐

