如何优化存在消费者空转问题的单生产者多消费者Java代码
生产者消费者模型优化方案
核心优化思路是借助ReentrantLock配套的Condition等待通知机制,替代原有空转轮询逻辑,彻底解决消费者空转浪费CPU的问题。
具体修改点
- 给
Assembly类新增绑定到锁的Condition实例,作为缓冲区非空的通知载体 - 消费者检测到缓冲区为空时,调用
await()方法主动挂起并释放锁,不再循环空转 - 生产者每次往缓冲区添加元素后,主动唤醒等待的消费者
- 处理终止标记99时额外唤醒所有消费者,避免剩余消费者永久挂起
- 空缓冲区判断改用
while循环,规避线程虚假唤醒问题
修改后完整代码
import java.util.List; import java.util.Random; import java.util.ArrayList; import java.util.concurrent.locks.Condition; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; public class Test { public static void main(String[] args) { Assembly assembly = new Assembly(new ArrayList<>(), new ReentrantLock(true)); new Thread(() -> assembly.consume()).start(); new Thread(() -> assembly.produce()).start(); new Thread(() -> assembly.consume()).start(); } } class Assembly { List<Integer> buffer; Lock bufferLock; // 新增非空等待条件 private final Condition notEmpty; public Assembly(List<Integer> buffer, Lock bufferLock) { this.buffer = buffer; this.bufferLock = bufferLock; this.notEmpty = bufferLock.newCondition(); } public void produce() { Integer[] nums = {1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 99}; Random random = new Random(); for (Integer num : nums) { try { bufferLock.lock(); buffer.add(num); // 新增元素后唤醒一个等待的消费者 notEmpty.signal(); if (num != 99) { System.out.println("Added: " + num); } } finally { bufferLock.unlock(); try { Thread.sleep(random.nextInt(1000)); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } } } public void consume() { while (true) { try { bufferLock.lock(); // 用while循环判断,避免虚假唤醒 while (buffer.isEmpty()) { notEmpty.await(); } if (buffer.get(0).equals(99)) { // 唤醒所有等待的消费者,避免永久挂起 notEmpty.signalAll(); break; } System.out.println("Removed: " + buffer.remove(0) + " by " + Thread.currentThread().getName()); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } finally { bufferLock.unlock(); } } } }
补充说明
- 单次新增元素用
signal()比signalAll()更高效,不会产生惊群效应,仅唤醒一个消费者处理即可 - 如果后续需要限制缓冲区最大长度,可以新增另一个
notFull条件,控制生产者在缓冲区满时挂起等待 - 所有抛出
InterruptedException的位置都主动恢复了线程中断状态,符合并发编程最佳实践
内容的提问来源于stack exchange,提问作者Rajat Aggarwal
相关产品推荐
相关产品推荐

