【发布时间】:2020-11-09 16:44:55
【问题描述】:
我在应用程序中有原始的消息传递系统。消息可以由生产者从一个线程提交并由消费者在另一个线程中处理 - 设计只有两个线程:一个线程用于消费者,另一个线程用于生产者,这是不可能的改变这个逻辑。
我正在使用ConcurrentLinkedQueue<> 实现来处理消息:
// producer's code (adds the request)
this.queue.add(req);
// consumer's code (busy loop with request polling)
while (true) {
Request req = this.queue.poll();
if (req == null) {
continue;
}
if (req.last()) {
// last request submitted by consumer
return;
}
// function to process the request
this.process(req);
}
处理逻辑非常快,消费者每秒可能会收到大约X_000_000 个请求。
但我发现使用 profiler 时 queue.poll() 有时非常慢(似乎是当队列从生产者那里接收大量新项目时) - 与已经填充的相比,接收大量新消息时速度慢约 10 倍在不从另一个线程添加新项目的情况下排队。
可以优化吗?对于这种特殊情况,最好的Queue<> 实现是什么(poll() 一个线程,add() 一个线程)?也许自己实现一些简单的队列会更容易?
【问题讨论】:
-
这能回答你的问题吗? LinkedBlockingQueue vs ConcurrentLinkedQueue
-
@Amongalen 感谢您的链接,但我不这么认为 -
LinkedBlockingQueue在我的情况下性能更差:我尝试使用它,但它比非慢 2 倍以上-阻止ConcurrentQueue。所以这两种队列实现都不适合我的应用。 -
你试过 SynchronousQueue 吗?当一个线程想要将数据传递给另一个线程时使用此类。
-
@Eric 生产者如何能够在消费者忙于处理一个元素时将多个元素添加到队列中以存储它们?
-
对于它的价值,我过去已经成功地在这种情况下使用了循环缓冲区,而且工作量相对较小。如果我没记错的话,只需要一个(包装)写入索引、一个(包装)读取索引和一个监视器变量,供写入器启动等待的读取器。读取器从不更新写入索引,写入器从不更新读取索引。 - 我记得,我使用的是预先分配的固定长度数组。动态调整队列大小可能会更复杂。
标签: java multithreading algorithm concurrency queue