You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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()的线程,没有持有方法所属对象的监视器锁。当前代码的并发逻辑存在多处核心错误:

  1. 锁对象不匹配
    run()方法中synchronized(this)持有的是Producer/Consumer实例自身的锁,但代码中调用producerMutex.WorkingQueue.notifyAll()、consumerMutex.WorkingQueue.notifyAll()时,从未获取过WorkingQueue对象的锁,直接调用notify方法必然抛出异常。此外无参wait()是在Producer/Consumer实例上等待,生产者和消费者持有的是完全独立的锁对象,根本不可能实现跨线程唤醒。
  2. 锁释放逻辑错误
    缓冲区为空/满的分支中,代码先调用wait()让线程进入休眠,之后才写了mutex和semaphore的release逻辑。线程进入等待状态后不会执行后续代码,锁会被永久占用,其他线程永远无法获取锁访问临界区。
  3. 缺少循环状态校验
    等待逻辑没有放在循环中,线程被唤醒后不会重新校验缓冲区状态,会出现虚假唤醒问题:比如消费者被唤醒时缓冲区仍为空,直接执行pop操作会引发数据错误。
  4. 临界区未被完全保护
    缓冲区空/满判断、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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.02 05:15:18