基于BlockingQueue的生产者/消费者复合操作同步难题
基于BlockingQueue的生产者/消费者模式复合操作原子性问题
1. 问题概述
我在实现基于BlockingQueue的生产者/消费者模式时,遇到了复合操作的原子性问题,大概率是忽略了某些关键细节。我的核心需求是:
- 消费者的「从队列取出对象 + 对该对象执行后续消费操作」序列必须具备原子性
- 生产者的「将对象放入队列 + 对该对象执行后续生产操作」序列必须具备原子性
- 上述两个原子操作序列要基于同一对象同步
如果不保证这种原子性,会引发问题——比如生产者代码里注释标记的「PROBLEM!!」场景。但我不能直接在take()调用和关联消费操作外层加synchronized块,因为队列空时,消费者会持有同步锁永久阻塞,直接阻止生产者进入临界区执行生产操作。
2. 简化示例代码
公共代码
Queue<QObj> nbq = new ConcurrentLinkedQueue(); BlockingQueue<QObj> bq = new LinkedBlockingQueue<>(); List<String> idList = new LinkedList<>(); Object lockObj = idList; int Idx = 1; public static class QObj { public String id; public String content; public QObj(String id, String content) { this.id = id; this.content = content; } }
生产者核心逻辑
public void produceBlocking() { QObj o = new QObj(String.valueOf(Idx), "Content_" + Idx++); // synchronized(lockObj) { // 把Queue.offer(...)放进同步块没意义,因为消费者那边没法用synchronized // 原因前面已经说过了 bq.offer(o); synchronized (lockObj) { // PROBLEM!! 到这一步,'o'可能已经被消费了 // 所以下面的操作不应该执行: // 执行生产者复合操作的关联部分 idList.add(o.id); // 执行该复合操作的其他步骤... } // } }
消费者核心逻辑
public void consumeBlocking() { while (true) { try { // synchronized (lockObj) { // 不能直接在这加synchronized让后续复合操作原子化 // - 队列空时,消费者会一直卡在这,因为还持有lockObj, // 会阻止生产者进入临界区生产 QObj o = bq.take(); synchronized (lockObj) { // 执行消费者复合操作的关联部分 idList.remove(o.id); // 执行该复合操作的其他步骤... } // } } catch (InterruptedException e) { e.printStackTrace(); } } }
3. 为何这不是普遍问题?
我觉得用BlockingQueue时这应该是个常见问题,但找不到直接的解决方案,这让我怀疑自己的理解有根本性错误。希望有人能给出直接的解决方案,或者指出我的理解误区。
4. 备选思路
我想了几种备选方案,但都没法直接解决问题,还存在缺陷(代码注释里标记的「DRAWBACK!!」):
4.1 - 执行前用Queue.contains()检查
public void produceBlockingWithCheck() { QObj o = new QObj(String.valueOf(Idx), "Content_" + Idx++); bq.offer(o); synchronized (lockObj) { // 先检查对象是否已经被消费 // DRAWBACK!! 这可能非常耗时,比如: // 当'bq'是LinkedBlockingQueue时,contains(...)会触发顺序遍历, // 队列很大时性能极差 if (bq.contains(o)) { // 执行生产者复合操作的关联部分 idList.add(o.id); // 执行该复合操作的其他步骤... } } }
4.2 - 调整生产者操作顺序,把Queue.offer()移到最后
public void produceBlockingOrderAdjusted() { QObj o = new QObj(String.valueOf(Idx), "Content_" + Idx++); // 先执行生产者复合操作的关联部分,再调用BlockingQueue.offer(...) // DRAWBACK!! 就算这个简单场景能工作,不是所有场景都能调整顺序吧? synchronized (lockObj) { idList.add(o.id); // 执行该复合操作的其他步骤... } bq.offer(o); }
4.3 - 改用非阻塞队列
public void produceNonBlocking() { QObj o = new QObj(String.valueOf(Idx), "Content_" + Idx++); synchronized(lockObj) { nbq.offer(o); // 执行生产者复合操作的关联部分 idList.add(o.id); // 执行该复合操作的其他步骤... } } public void consumeNonBlocking() { while (true) { synchronized (lockObj) { // 自己实现阻塞逻辑 QObj o = nbq.poll(); if (o != null) { // 执行消费者复合操作的关联部分 idList.add(o.id); // 执行该复合操作的其他步骤... } // DRAWBACK!! 如果生产者生产速度赶不上消费者消费速度, // 这种空轮询会频繁发生,非常耗费资源 } } }
内容的提问来源于stack exchange,提问作者coder joe
相关产品推荐
相关产品推荐

