【发布时间】: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