【问题标题】:How to take items from queue in chunks?如何从队列中分块取出项目?
【发布时间】:2018-05-17 13:21:59
【问题描述】:

我有多个生产者线程同时将对象添加到共享队列。

我想创建一个单线程使用者,从该共享队列中读取数据以进行进一步的数据处理(数据库批量插入)。

问题:我只想从队列中分块获取数据,以便在批量插入期间获得更好的性能。因此,我必须以某种方式检测队列中有多少项目,然后从队列中取出所有这些项目,然后再次清空队列。

 BlockingQueue<Integer> sharedQueue = new LinkedBlockingQueue<>();

 ExecutorService pes = Executors.newFixedThreadPool(4);
 ExecutorService ces = Executors.newFixedThreadPool(1);

 pes.submit(new Producer(sharedQueue, 1));
 pes.submit(new Producer(sharedQueue, 2));
 pes.submit(new Producer(sharedQueue, 3));
 pes.submit(new Producer(sharedQueue, 4));
 ces.submit(new Consumer(sharedQueue, 1));

class Producer implements Runnable {
    run() {
            ...
            sharedQueue.put(obj);
    }
}

class Consumer implements Runnable {
    run() {
            ...
            sharedQueue.take();
    }
}

消费者的问题:我如何轮询共享队列,等待队列有 X 个项目,然后取出所有项目并同时清空队列(以便消费者可以重新开始轮询和等待)?

我愿意接受任何建议,不一定受上述代码的约束。

【问题讨论】:

  • 你总是只有一个消费者吗?
  • 是的,消费者总是单线程的。
  • Guava 的 Queues.drain() 就是这样做的。

标签: java concurrency queue


【解决方案1】:

我最近开发了这个实用程序,如果队列元素没有达到批量大小,它会使用刷新超时来批量处理 BlockingQueue 元素。 它还支持使用多个实例来详细说明同一组数据的扇出模式:

// Instantiate the registry
FQueueRegistry registry = new FQueueRegistry();

// Build FQueue consumer
registry.buildFQueue(String.class)
                .batch()
                .withChunkSize(5)
                .withFlushTimeout(1)
                .withFlushTimeUnit(TimeUnit.SECONDS)
                .done()
                .consume(() -> (broadcaster, elms) -> System.out.println("elms batched are: "+elms.size()));

// Push data into queue
for(int i = 0; i < 10; i++){
        registry.sendBroadcast("Sample"+i);
}

更多信息在这里!

https://github.com/fulmicotone/io.fulmicotone.fqueue

【讨论】:

    【解决方案2】:

    您最好在消费者中创建一个内部List 并从队列中获取对象并添加到该列表中,而不是检查队列的大小。一旦列表有 X 项,您就可以进行处理,然后清空内部列表。

    class Consumer implements Runnable {
      private List itemsToProcess = new ArrayList();
      run() {
            while (true) { // or when producers are stopped and queue is empty
              while (itemsToProcess.size() < X) {
                itemsToProcess.add(sharedQueue.take());
              }
              // do process
              itemsToProcess.clear();
            }
      }
    }
    

    您可以使用 BlockingQueue.poll(timeout) 而不是 BlockingQueue.take() 和一些合理的超时并检查 null 的结果,以检测所有生产者都完成并且队列为空以便能够关闭您的消费者的情况。

    【讨论】:

    • 这是一个很好的建议。在我的情况下,一个简单的带有 nullcheck 的 queue.poll() 就可以了。
    猜你喜欢
    • 1970-01-01
    • 2014-07-26
    • 1970-01-01
    • 1970-01-01
    • 2015-06-01
    • 1970-01-01
    • 2010-09-14
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多