【问题标题】:Wait for all blocking queue elements to be processed after they are taken out等待所有阻塞队列元素取出后处理
【发布时间】:2018-08-14 19:14:31
【问题描述】:

在以下场景中,终结器线程必须等待消费者线程处理完所有队列元素才能完成执行:

private final BlockingQueue<Object> queue = new LinkedBlockingQueue<>();
private final Object queueMonitor = new Object();

// Consumer thread
while (true) {
    Object element = queue.take();
    consume(element);
    synchronized (queueMonitor) {
        queueMonitor.notifyAll();
    }
}

// Finalizer thread
synchronized (queueMonitor) {
    while (!queue.isEmpty()) {
        queueMonitor.wait();
    }
}

元素会随着时间的推移添加到队列中。 消费者守护线程一直运行,直到 JVM 终止,此时必须允许它完成对所有排队元素的处理。 目前这是由终结器线程完成的,它是一个关闭钩子,应该延迟在 JVM 终止时杀死消费者线程。

问题:
如果在从队列中取出最后一个元素后启动终结器线程,则 while 循环条件的计算结果为 false,因此在 consume() 尚未返回时执行完成,因为完全跳过了等待 queueMonitor

研究:
一个理想的解决方案是peek the queue,然后在元素被消耗后将其删除。

【问题讨论】:

  • 您是说您希望“终结器”线程在“消费者”线程完成之前什么都不做?在这种情况下,为什么不在一个线程中完成这两项工作呢?
  • 如果你的“消费者”线程不是守护进程怎么办?如果它不是while(true),而是循环直到在队列中找到一颗毒丸,然后退出呢?然后,无论您调用什么函数来关闭应用程序,它都可以将毒丸送入队列,然后join 消费者。
  • @jameslarge 没有调用函数来关闭应用程序。该代码是库的一部分,我不希望客户端调用 API 来中断非守护线程并让 JVM 终止,尽管这是一个选项。
  • 有什么理由不将消费者中的整个 while 循环包装在同步块中锁定 queueMonitor 以使终结器块直到消费者完成?我想说整个设置看起来很容易出错,我建议使用CountDownLatch 之类的东西,如下所示。
  • @JanusVarmarken 原因是消费者线程永远不会完成。它是一个在 JVM 终止时被杀死的守护线程。

标签: java multithreading concurrency blockingqueue


【解决方案1】:

一种方法可能是您使用CountDownLatch - 在其上放置终结器块,并在consume() 之后使用消费者调用倒计时。

基本上不在队列上阻塞,在任务完成时阻塞。

private final BlockingQueue<Object> queue = new LinkedBlockingQueue<>();
private volatile boolean running = true;
private final CountDownLatch terminationLatch = new CountDownLatch(1);

// Consumer thread
while (running || !queue.isEmpty()) {
    Object element = queue.poll(100, TimeUnit.MILLISECONDS);
    if (element == null) continue;
    consume(element);
}
terminationLatch.countDown();

// Finalizer thread
running = false;
terminationLatch.await();

【讨论】:

  • 重构以使用 ExecutorService,终结器可以调用关闭并阻塞 awaitTermination - 没有必要重新发明轮子。
  • 不幸的是,这是不可能的。 consume() 的实现会安排额外的任务,并会抛出 RejectedExecutionException
  • 好的,在这种情况下,我只需使用 CountDownLatch,将其称为 exitLatch 并使用计数 1 进行初始化 - 消费者线程运行 while(queue.size() &gt; 0) 然后在一段时间之后循环完成,它调用exitLatch.countDown()。同时,终结器线程所要做的就是exitLatch.await(1_000,TimeUnit.MILLISECONDS)(+ 处理中断的异常等)
  • 元素会随着时间的推移添加到队列中。消费者线程一直运行,直到 JVM 终止,此时必须允许它完成所有剩余元素的处理。
  • 好的,那么请保留您的while(true) - 其他一切仍然适用
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2019-10-16
  • 2016-06-03
  • 1970-01-01
  • 2014-10-16
  • 2013-02-21
  • 1970-01-01
  • 2019-01-20
相关资源
最近更新 更多