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

多生产者消费者有界缓冲区问题:生产者结束后消费者无限等待解决方案

多生产者消费者有界缓冲区程序无法终止的解决方案

问题核心是生产者完成生产后,消费者没有收到“生产结束”的信号,仍在等待新元素,导致线程挂起、程序无法退出。下面是可行的解决思路和实现代码:

核心思路

给缓冲区添加生产结束标记和活跃生产者计数器:

  • 用计数器跟踪当前还在生产的生产者数量,所有生产者完成后,标记“生产结束”。
  • 消费者在等待时,同时判断缓冲区是否为空、生产是否结束;如果生产结束且缓冲区空,直接退出线程。

完整实现代码

缓冲区类(核心修改)

class BoundedBuffer {
    private final Object[] buffer;
    private int in, out, count;
    // 标记所有生产者是否已完成生产
    private volatile boolean productionDone = false;
    // 跟踪当前活跃的生产者数量
    private volatile int activeProducers;

    public BoundedBuffer(int capacity) {
        buffer = new Object[capacity];
        activeProducers = 0;
    }

    // 生产者启动时调用,注册为活跃生产者
    public synchronized void registerProducer() {
        activeProducers++;
    }

    public synchronized void put(Object item) throws InterruptedException {
        while (count == buffer.length) {
            wait();
        }
        buffer[in] = item;
        in = (in + 1) % buffer.length;
        count++;
        notifyAll();
    }

    public synchronized Object take() throws InterruptedException {
        // 缓冲区为空且生产未结束时,继续等待
        while (count == 0 && !productionDone) {
            wait();
        }
        // 生产结束且缓冲区空,返回null表示无更多元素
        if (count == 0 && productionDone) {
            return null;
        }
        Object item = buffer[out];
        out = (out + 1) % buffer.length;
        count--;
        notifyAll();
        return item;
    }

    // 生产者完成生产时调用
    public synchronized void finishProduction() {
        activeProducers--;
        // 所有生产者都完成,标记生产结束并唤醒所有等待的消费者
        if (activeProducers == 0) {
            productionDone = true;
            notifyAll();
        }
    }
}

生产者线程

class Producer implements Runnable {
    private final BoundedBuffer buffer;
    private final int itemsToProduce;

    public Producer(BoundedBuffer buffer, int itemsToProduce) {
        this.buffer = buffer;
        this.itemsToProduce = itemsToProduce;
    }

    @Override
    public void run() {
        buffer.registerProducer();
        try {
            for (int i = 0; i < itemsToProduce; i++) {
                Object item = "Item " + Thread.currentThread().getId() + "-" + i;
                buffer.put(item);
                System.out.println("Produced: " + item);
                Thread.sleep(100); // 模拟生产耗时
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        } finally {
            // 无论是否中断,都标记当前生产者完成
            buffer.finishProduction();
        }
    }
}

消费者线程

class Consumer implements Runnable {
    private final BoundedBuffer buffer;

    public Consumer(BoundedBuffer buffer) {
        this.buffer = buffer;
    }

    @Override
    public void run() {
        try {
            Object item;
            // 拿到null表示生产结束且无剩余元素,退出循环
            while ((item = buffer.take()) != null) {
                System.out.println("Consumed by " + Thread.currentThread().getId() + ": " + item);
                Thread.sleep(200); // 模拟消费耗时
            }
            System.out.println("Consumer " + Thread.currentThread().getId() + " exited: no more items.");
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

测试主类

public class Main {
    public static void main(String[] args) {
        BoundedBuffer buffer = new BoundedBuffer(5);
        // 2个生产者,每个生产3个元素
        new Thread(new Producer(buffer, 3)).start();
        new Thread(new Producer(buffer, 3)).start();
        // 2个消费者
        new Thread(new Consumer(buffer)).start();
        new Thread(new Consumer(buffer)).start();
    }
}

关键注意事项

  • volatile修饰productionDone和activeProducers,确保线程间状态可见,避免读取到过期值。
  • 必须用notifyAll()唤醒线程,不能用notify()——否则可能只有部分消费者收到通知,剩余的继续阻塞。
  • 生产者在finally块中调用finishProduction(),保证即使生产过程被中断,也能正确更新状态。

内容的提问来源于stack exchange,提问作者HM Try

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 14:05:23