【问题标题】:Suspend consumer in producer/consumer pattern在生产者/消费者模式中暂停消费者
【发布时间】:2015-10-12 09:23:12
【问题描述】:

我有生产者和消费者与BlockingQueue 连接。

消费者等待队列中的记录并进行处理:

Record r = mQueue.take();
process(r);

我需要从其他线程暂停这个过程一段时间。如何实现?

现在我想这样实现它,但它看起来像一个糟糕的解决方案:

private Object mLock = new Object();
private boolean mLocked = false;

public void lock() {
    mLocked = true;
}

public void unlock() {
    mLocked = false;
    mLock.notify();

}

public void run() {
    ....
            Record r = mQueue.take();
            if (mLocked) {
                mLock.wait();
            }
            process(r);
}

【问题讨论】:

  • 您介意提供更多关于您的特定需求的背景信息吗?我问的原因是,如果我真的需要暂停链中的某个人,我宁愿暂停生产者而不是消费者。它基本上实现了类似的效果,而不会在暂停生效时队列增长到容量。

标签: java multithreading producer-consumer blockingqueue


【解决方案1】:

我认为您的解决方案简单而优雅,并且认为您应该对其进行一些修改。我提出的修改是synchronization

没有它,线程干扰和内存一致性错误可能(并且经常发生)发生。最重要的是,您不能在您不拥有的锁上使用waitnotify(如果您在synchronized 块中拥有它,则您拥有它。)。修复很简单,只需在您等待/通知的位置添加一个mLock 同步块。此外,当您从不同的线程更改mLocked 时,您需要将其标记为volatile

private Object mLock = new Object();
private volatile boolean mLocked = false;

public void lock() {
    mLocked = true;
}

public void unlock() {
    synchronized(mlock) {
        mLocked = false;
        mLock.notify();
    }

}

public void run() {
    ....
            Record r = mQueue.take();
            synchronized(mLock) {
                while (mLocked) {
                    mLock.wait();
                }
            }
            process(r);
}

【讨论】:

【解决方案2】:

您可以根据相同的条件使用java.util.concurrent.locks.ConditionJava docs暂停一段时间。

这种方法对我来说看起来很干净,ReentrantLock 机制的吞吐量比同步的要好。阅读以下摘自IBM article

作为奖励,ReentrantLock 的实现更具可扩展性 在竞争下比当前执行同步。 (它 竞争性能可能会有所改善 在 JVM 的未来版本中同步。)这意味着当 许多线程都在争夺同一个锁,总数 ReentrantLock 的吞吐量通常会比 ReentrantLock 更好 同步。


BlockingQueue 以解决生产者-消费者问题而闻名,它也使用Condition 等待。

请参阅以下示例,取自 Java 文档的 Condition,这是生产者 - 消费者模式的示例实现。

class BoundedBuffer {
   final Lock lock = new ReentrantLock();
   final Condition notFull  = lock.newCondition(); 
   final Condition notEmpty = lock.newCondition(); 

   final Object[] items = new Object[100];
   int putptr, takeptr, count;

   public void put(Object x) throws InterruptedException {
     lock.lock();
     try {
       while (count == items.length)
         notFull.await();
       items[putptr] = x;
       if (++putptr == items.length) putptr = 0;
       ++count;
       notEmpty.signal();
     } finally {
       lock.unlock();
     }
   }

   public Object take() throws InterruptedException {
     lock.lock();
     try {
       while (count == 0)
         notEmpty.await();
       Object x = items[takeptr];
       if (++takeptr == items.length) takeptr = 0;
       --count;
       notFull.signal();
       return x;
     } finally {
       lock.unlock();
     }
   }
 }

进一步阅读:

【讨论】:

  • 先生。 down 投票者,你能解释一下你的投票吗?
  • 当这个人实际上需要一些时间来提供有见地的解释时,我讨厌投反对票。如果此解决方案对您不起作用,那很好。但它不值得一票否决。 +1 对我来说,即使我发现第一个更简单,第二个提供关于线程和并发性的见解。
  • 不是我的反对意见,但根据 10 多年前的文章给出性能判断似乎不太可靠。
  • @PatB 感谢朋友的支持和支持。
  • @NathanHughes 我会接受您的评论,我试图获得一些最新的参考资料,但在我的快速搜索中找不到任何参考资料。但我认为这将成立,因为这与 BlockingQueue 实现中使用的机制相同,如 ArrayBlockingQueueLinkedBlockingQueue
【解决方案3】:

创建一个扩展BlockingQueue 实现的新类。添加两个新方法pause()unpause()。需要时,考虑paused 标志并使用另一个blockingQueue2 等待(在我的示例中仅在take() 方法中,而不在put() 中):

public class BlockingQueueWithPause<E> extends LinkedBlockingQueue<E> {

    private static final long serialVersionUID = 184661285402L;

    private Object lock1 = new Object();//used in pause() and in take()
    private Object lock2 = new Object();//used in pause() and unpause()

    //@GuardedBy("lock")
    private volatile boolean paused;

    private LinkedBlockingQueue<Object> blockingQueue2 = new LinkedBlockingQueue<Object>();

    public void pause() {
        if (!paused) {
            synchronized (lock1) {
            synchronized (lock2) {
                if (!paused) {
                    paused = true;
                    blockingQueue2.removeAll();//make sure it is empty, e.g after successive calls to pause() and unpause() without any consumers it will remain unempty
                }
            }
            }
        }
    }

    public void unpause() throws InterruptedException {
        if (paused) {
            synchronized (lock2) {
                paused = false;
                blockingQueue2.put(new Object());//release waiting thread, if there is one
            }
        }
    }

    @Override
    public E take() throws InterruptedException {
        E result = super.take();

        if (paused) {
            synchronized (lock1) {//this guarantees that a single thread will be in the synchronized block, all other threads will be waiting
                if (paused) {
                    blockingQueue2.take();
                }
            }
        }

        return result;
    }

    //TODO override similarly the poll() method.
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-10-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多