【问题标题】:Producer-consumer problem with a twist扭曲的生产者消费者问题
【发布时间】:2011-04-30 15:12:32
【问题描述】:
  • 生产者是有限的,消费者也应该如此。
    • 问题是何时停止,而不是如何运行。
  • 可以通过任何 类型的 BlockingQueue 进行通信。
    • 不能依赖中毒队列(PriorityBlockingQueue)
    • 不能依赖锁定队列(SynchronousQueue)
    • 不能只依赖offer/poll(SynchronousQueue)
    • 可能存在更奇特的队列。

在另一个(可能是惰性的)seq 上创建一个排队的 seq。排队的 seq 会在后台产生一个具体的 seq,最多可以达到 n 项领先于消费者。 n-or-q 可以是整数 n 缓冲区 大小,或 java.util.concurrent BlockingQueue 的一个实例。笔记 如果读者领先于 制片人。

http://clojure.github.com/clojure/clojure.core-api.html#clojure.core/seque

到目前为止我的尝试 + 一些测试:https://gist.github.com/934781

感谢 Java 或 Clojure 中的解决方案。

【问题讨论】:

  • 在 Java 中,我只会使用 ExecutorService 或包装它的类来处理您需要的所有事件类型,因为它可以处理您提到的所有事件等等。

标签: java concurrency clojure


【解决方案1】:
class Reader {

    private final ExecutorService ex = Executors.newSingleThreadExecutor();
    private final List<Object> completed = new ArrayList<Object>();
    private final BlockingQueue<Object> doneQueue = new LinkedBlockingQueue<Object>();
    private int pending = 0;

    public synchronized Object take() {
        removeDone();
        queue();
        Object rVal;
        if(completed.isEmpty()) {
            try {
                rVal = doneQueue.take();
            } catch (InterruptedException e) {
                throw new RuntimeException(e);
            }
            pending--;
        } else {
            rVal = completed.remove(0);
        }
        queue();
        return rVal;
    }

    private void removeDone() {
        Object current = doneQueue.poll();
        while(current != null) {
            completed.add(current);
            pending--;
            current = doneQueue.poll();
        }
    }

    private void queue() {
        while(pending < 10) {
            pending++;
            ex.submit(new Runnable() {

                @Override
                public void run() {
                    doneQueue.add(compute());
                }

                private Object compute() {
                    //do actual computation here
                    return new Object();
                }
            });
        }
    }
}

【讨论】:

  • 1个后台线程最多提前产生10个结果,由于前台取结果,后台线程产生新结果,如果没有结果,取块直到有新结果为止跨度>
  • 嗯,是的,但我说生产者是有限的,主要问题是停止。如果没有更多结果,这将永远坐在那里。
  • 当您说“不能依赖毒化队列(PriorityBlockingQueue)”时,您的意思是您不能向 doneQueue 添加 DONE 标记,或者您不能使用 PriorityBlockingQueue?
  • 我的意思是当您插入标记时,PriorityBlockingQueue 会在内部深处引发异常。
  • 应该是可以解决的,可能是因为你插入了null,或者你插入了一个与队列的其他元素不兼容的DONE标记
【解决方案2】:

恐怕不完全是答案,而是一些评论和更多问题。我的第一个答案是:使用clojure.core/seque。生产者需要以某种方式传达序列结束,以便消费者知道何时停止,并且我假设生产元素的数量事先不知道。为什么不能使用 EOS 标记(如果这就是队列中毒的意思)?

如果我正确理解您的替代 seque 实现,当元素从您的函数之外的队列中取出时,它将中断,因为在这种情况下 channelq 将不同步:通道将容纳更多 @ 987654325@ 元素比q 中的元素多,导致它阻塞。可能有一些方法可以确保 channelq 始终保持同步,但这可能需要实现您自己的 Queue 类,而且它增加了如此多的复杂性,我怀疑它是否值得。

此外,您的实现不区分正常 EOS 和由于线程中断导致的异常队列终止 - 取决于您使用它的原因,您可能想知道哪个是哪个。就我个人而言,我不喜欢以这种方式使用异常——将异常用于异常情况,而不是用于正常的流控制。

【讨论】:

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