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

多线程新手求助:生产者-消费者问题中双方均陷入等待排查

嘿,我来帮你排查这个生产者-消费者模型里的问题,你遇到的双方都陷入等待的情况其实是典型的线程间通信错误导致的死锁,咱们一步步拆解来看:

你的代码中存在的核心问题

1. 条件判断用if而非while,踩了虚假唤醒的坑

Java的wait()方法可能会出现虚假唤醒——线程没被notify也会自己醒过来,就算是被正常唤醒,当前的buffer状态也可能已经变化(比如多线程场景下其他线程抢先修改了buffer)。你原代码里用if判断buffer状态:

// Producer类里的判断
if (bf.isFull()) { ... }
// Consumer类里的判断
if (bf.isEmpty()) { ... }

这会导致线程被唤醒后直接往下执行,不会重新检查buffer状态,很容易引发错误或者死锁。正确姿势是用while循环持续校验条件。

2. 生产/消费后没及时通知对方

你的代码只有在自己要进入等待状态时才通知对方,但正常生产或消费完,完全没管对方是不是在等着。举个例子:

  • 消费者消费了一个元素,buffer从满变空?不,是从满变不满,这时候生产者可能还在等空位,但你没通知它继续生产;
  • 生产者生产了一个元素,buffer从空变非空,消费者可能还在等元素,你也没通知它消费。
    这就直接导致对方线程一直卡在等待里,最后双方都动不了。

3. Buffer类的设计不符合生产者-消费者逻辑

你的Buffer.add()在buffer满时直接抛异常,正常的生产者应该是等buffer有空位再生产,而不是直接崩溃。而且把等待/通知的逻辑丢给Producer和Consumer处理,也容易搞混锁对象。

4. 线程终止逻辑没处理等待状态

terminate()只把running设为false,但如果线程正卡在wait()里,根本不会察觉到这个标志,会一直堵在那里退不出。得在终止时唤醒等待的线程才行。


修复后的完整代码

我把所有问题都修正了,关键位置加了注释:

Main类

package producer.consumer2;
import java.util.Scanner;
public class Main {
    public static void main(String[] args) {
        Buffer<Integer> bf = new Buffer<>(10);
        Producer prod = new Producer(bf);
        Consumer cons = new Consumer(bf);
        
        new Thread(prod).start();
        new Thread(cons).start();
        
        if(quitInput()) {
            prod.terminate();
            cons.terminate();
            // 唤醒所有等在Buffer上的线程,确保能正常退出
            synchronized (bf) {
                bf.notifyAll();
            }
        }
    }
    private static boolean quitInput() {
        Scanner sc = new Scanner(System.in);
        String line;
        do {
            line = sc.nextLine();
            if(line.toLowerCase().equals("q") || line.toLowerCase().equals("quit")) {
                sc.close();
                return true;
            }
        } while(true);
    }
}

Buffer类(核心优化)

package producer.consumer2;
import java.util.ArrayList;
public class Buffer<E> {
    private final int MAX_LENGTH;
    private ArrayList<E> values;
    public Buffer(int length){
        MAX_LENGTH = length;
        values = new ArrayList<>(length);
    }
    
    // 生产者调用:满了就等,生产完通知消费者
    public synchronized void add(E e) throws InterruptedException {
        // while循环防虚假唤醒,确保条件满足才执行
        while(values.size() >= MAX_LENGTH) {
            System.out.println("Buffer满了,生产者等待中...");
            wait();
        }
        values.add(e);
        System.out.println("生产了: " + e + " | 当前Buffer: " + values);
        notifyAll(); // 通知消费者可以消费了
    }
    
    // 消费者调用:空了就等,消费完通知生产者
    public synchronized E remove() throws InterruptedException {
        while(values.isEmpty()) {
            System.out.println("Buffer空了,消费者等待中...");
            wait();
        }
        E val = values.remove(0);
        System.out.println("消费了: " + val + " | 当前Buffer: " + values);
        notifyAll(); // 通知生产者可以生产了
        return val;
    }
}

Consumer类

package producer.consumer2;
public class Consumer implements Runnable {
    private final Buffer<Integer> bf;
    private volatile boolean running = true;
    public Consumer(Buffer<Integer> bf) {
        this.bf = bf;
    }
    @Override
    public void run() {
        int sum = 0;
        int counter = 0;
        while (running) {
            try {
                int val = bf.remove();
                sum += val;
                counter++;
                // 加个小延迟模拟消费耗时,方便观察流程
                Thread.sleep(200);
            } catch (InterruptedException e) {
                // 收到中断信号,重置状态后退出循环
                Thread.currentThread().interrupt();
                break;
            }
        }
        // 避免counter为0时除以0的情况
        if (counter > 0) {
            System.out.println("累计消费" + counter + "个数字,平均值: " + ((double) sum / counter));
        } else {
            System.out.println("没有消费任何数字");
        }
    }
    public void terminate() {
        this.running = false;
        // 中断线程,让卡在wait的线程抛出异常退出
        Thread.currentThread().interrupt();
    }
}

Producer类

package producer.consumer2;
public class Producer implements Runnable {
    private final Buffer<Integer> bf;
    private volatile boolean running = true;
    public Producer(Buffer<Integer> bf) {
        this.bf = bf;
    }
    @Override
    public void run() {
        int counter = 1;
        while (running) {
            try {
                bf.add(counter);
                counter++;
                // 加个小延迟模拟生产耗时
                Thread.sleep(100);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                break;
            }
        }
        System.out.println("生产者已终止");
    }
    public void terminate() {
        this.running = false;
        Thread.currentThread().interrupt();
    }
}

修复后的关键改进点
  1. Buffer封装所有线程安全逻辑:把等待/通知、生产消费都放在Buffer里,所有线程基于同一个Buffer锁通信,逻辑更清晰;
  2. while循环校验条件:彻底解决虚假唤醒和状态变化的问题;
  3. 生产/消费后立即notifyAll:确保对方线程能及时被唤醒,不会永久等待;
  4. 完善的终止逻辑:通过中断+标志位,无论线程是在运行还是等待,都能正常退出;
  5. 移除不必要的互相引用:Producer和Consumer不再需要持有对方对象,降低耦合度。

内容的提问来源于stack exchange,提问作者Čermák

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 09:09:42