多生产者消费者有界缓冲区问题:生产者结束后消费者无限等待解决方案
多生产者消费者有界缓冲区程序无法终止的解决方案
问题核心是生产者完成生产后,消费者没有收到“生产结束”的信号,仍在等待新元素,导致线程挂起、程序无法退出。下面是可行的解决思路和实现代码:
核心思路
给缓冲区添加生产结束标记和活跃生产者计数器:
- 用计数器跟踪当前还在生产的生产者数量,所有生产者完成后,标记“生产结束”。
- 消费者在等待时,同时判断缓冲区是否为空、生产是否结束;如果生产结束且缓冲区空,直接退出线程。
完整实现代码
缓冲区类(核心修改)
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
相关产品推荐
相关产品推荐

