在java中,有一个很好的包java.util.concurrent,它保存了BlockingQueue接口的实现.
我在Haskell中需要类似的东西,所以它能够
可能这可以通过STM或阻止事务来实现 - 但我无法在hackage上找到类似的东西.
我遇到了一个多线程服务器的问题,我正在构建一个学术练习,更具体地说是获得一个优雅地关闭的连接.
每个连接都由Session类管理.此类维护2个连接线程,一个DownstreamThread和一个UpstreamThread.
UpstreamThread在客户端套接字上阻塞,并将所有传入的字符串编码为消息,以传递到另一层进行处理.DownstreamThread阻塞BlockingQueue,其中插入了客户端的消息.当队列中有消息时,下游线程将消息从队列中取出,将其转换为字符串并将其发送到客户端.在最终系统中,应用程序层将对传入的消息进行操作,并将传出的消息推送到服务器以发送到适当的客户端,但是现在我只有一个简单的应用程序,在传回消息之前会休息一秒钟作为附加时间戳的传出消息.
我遇到的问题是当客户端断开连接时,整个事情都会正常关闭.我正在争论的第一个问题是正常的断开连接,客户端让服务器知道它正在结束与QUIT命令的连接.基本的伪代码是:
while (!quitting) {
inputString = socket.readLine () // blocks
if (inputString != "QUIT") {
// forward the message upstream
server.acceptMessage (inputString);
} else {
// Do cleanup
quitting = true;
socket.close ();
}
}
Run Code Online (Sandbox Code Playgroud)
上游线程的主循环查看输入字符串.如果是QUIT,则线程设置一个标志,表示客户端已结束通信并退出循环.这导致上游线程很好地关闭.
只要未设置连接关闭标志,下游线程的主循环就会等待BlockingQueue中的消息.如果是,则下游线程也应该终止.然而,它没有,它只是坐在那里等待.它的伪代码看起来像这样:
while (!quitting) {
outputMessage = messageQueue.take (); // blocks
sendMessageToClient (outputMessage);
}
Run Code Online (Sandbox Code Playgroud)
当我测试这个时,我注意到当客户端退出时,上游线程关闭,但下游线程没有.
经过一番搔痒之后,我意识到下游线程仍在阻塞BlockingQueue,等待永远不会到来的传入消息.上游线程不会在链上向前转发QUIT消息.
如何正常关闭下游线程?想到的第一个想法是在take()调用上设置超时.我不是太热衷于这个想法,因为无论你选择什么价值,它一定不会完全令人满意.要么它太长了,僵尸线程在那里停留了很长时间才关闭,或者它太短了,连接已经闲置了几分钟但仍然有效将会被杀死.我确实想过将QUIT消息发送到链中,但是这要求它完全往返服务器,然后是应用程序,然后再次返回服务器,最后返回到会话.这似乎也不是一个优雅的解决方案.
我确实查看了Thread.stop()的文档,但显然已经弃用了,因为它无论如何都没有正常工作,所以看起来它也不是一个真正的选项.我的另一个想法是强制在下游线程中以某种方式触发异常并让它在其最终块中清理,但这让我觉得这是一个可怕而又笨拙的想法.
我觉得两个线程都应该能够自己正常关闭,但我也怀疑如果一个线程结束它还必须发出信号通知另一个线程以更加主动的方式结束,而不是简单地为另一个线程设置一个标志来检查.由于我对Java仍然不是很有经验,所以我现在很缺乏想法.如果有人有任何建议,将不胜感激.
为了完整起见,我在下面列出了Session类的真实代码,不过我相信上面的伪代码片段涵盖了问题的相关部分.全班约250行.
import java.io.*;
import java.net.*;
import java.util.concurrent.*;
import java.util.logging.*;
/**
* Session class
*
* A session manages the individual connection between a …Run Code Online (Sandbox Code Playgroud) 如果队列已满,ArrayBlockingQueue将阻塞生产者线程,如果队列为空,它将阻塞使用者线程.
这种阻塞的概念是否违背了多线程的想法?如果我有一个'主'线程,让我们说我想将所有'Logging'活动委托给另一个线程.所以基本上在我的主线程中,我创建一个Runnable来记录输出,并将Runnable放在ArrayBlockingQueue上.这样做的全部目的是让'main'线程立即返回,而不会在昂贵的日志操作中浪费任何时间.
但是如果队列已满,主线程将被阻塞,并将等待一个点可用.那它对我们有什么帮助?
方法java.util.concurrent.BlockingQueue.add(E e)的JavaDoc读取:
布尔加法(E e)
如果可以在不违反容量限制的情况下立即执行此操作,则将指定的元素插入此队列,成功时返回true,如果当前没有可用空间则抛出IllegalStateException.使用容量限制队列时,通常最好使用offer.
我的问题是:它会不会返回虚假?如果没有,为什么这个方法返回一个布尔值?这对我来说似乎很奇怪.这背后的设计决策是什么?
谢谢你的知识!
曼努埃尔
对于我正在处理的日志记录功能,我需要有一个处理线程,当计数达到或超过一定数量时,它将等待作业并批量执行.由于这是生产者消费者问题的标准情况,我打算使用BlockingQueues.我有许多生产者使用add()方法向队列添加条目,而只有一个使用take()在队列中等待的消费者线程.
LinkedBlockingQueue似乎是一个不错的选择,因为它没有任何大小限制,但是我很困惑从文档中读取这个.
链接队列通常具有比基于阵列的队列更高的吞吐量,但在大多数并发应用程序中具有较低的可预测性能.
没有清楚地解释这个陈述的含义.有人可以点亮吗?这是否意味着LinkedBlockingQueue不是线程安全的?你们有没有遇到任何使用LinkedBlockingQueue的问题.
由于生产者的数量要多得多,因此总会有一种情况我可以遇到队列被大量条目淹没的情况.如果我改为使用ArrayBlockingQueue,它将队列的大小作为构造函数中的参数,我总是会遇到容量完全相关的异常.为了避免这种情况,我不知道如何确定我应该用实例化ArrayBlockingQueue的大小.您是否必须使用ArrayBlockingQueue解决类似的问题?
如果我们想要实现资源池,例如数据库连接池.您将使用哪个并发集合?BlockingQueue还是Semaphore?
因为BlockingQueue,就像生产者 - 消费者设计模式一样,生产者将所有连接放在队列上,而消费者将从队列中获取下一个可用连接.
对于Semaphore,您将信号量指定为池大小,并获取许可,直到达到池大小并等待其中任何一个释放许可并将资源放回池中.
哪一个更简单,更容易?什么是我们只能使用一个而不是其他的情况?
在使用a的应用程序中BlockingQueue,我面临的新要求只能通过迭代队列中存在的元素来实现(以提供有关元素当前状态的信息).
根据API JavadocBlockingQueue实现的唯一排队方法,只需要是线程安全的.Other API methods(例如,从中继承的那些Collection interface)可能不会同时使用,但我不确定这是否也适用于纯读取访问...
我可以安全地使用iterator() WITHOUT altering the producer/consumer threads哪些通常可以随时与队列交互?我不需要a 100% consistent iteration(在迭代队列时我是否看到添加/删除元素并不重要),但我不想最终讨厌ConcurrentModificationExceptions.
请注意,该应用程序当前正在使用a LinkedBlockingQueue,但我可以自由选择任何其他(unbounded) BlockingQueue implementation(包括free open-source third-party implementations).此外,我不想依赖将来可能会破坏的东西,所以我想要一个根据的解决方案,API并且不仅仅是恰好与之合作current JRE.
我需要在其队列已满时阻止ForkJoinPool上的线程.这可以在标准的ThreadPoolExecutor中完成,例如:
private static ExecutorService newFixedThreadPoolWithQueueSize(int nThreads, int queueSize) {
return new ThreadPoolExecutor(nThreads, nThreads,
5000L, TimeUnit.MILLISECONDS,
new ArrayBlockingQueue<Runnable>(queueSize, true), new ThreadPoolExecutor.CallerRunsPolicy());
}
Run Code Online (Sandbox Code Playgroud)
我知道,ForkJoinPool中有一些Dequeue,但是我无法通过它访问它.
更新:请参阅下面的答案.
的定义ExecutorService.newCachedThreadPool()是
public static ExecutorService newCachedThreadPool(ThreadFactory threadFactory) {
return new ThreadPoolExecutor(0, Integer.MAX_VALUE,
60L, TimeUnit.SECONDS,
new SynchronousQueue<Runnable>(),
threadFactory);
}
Run Code Online (Sandbox Code Playgroud)
它正在创建一个带有 的池corePoolSize = 0,maximumPoolSize = Integer.MAX_VALUE以及一个无界队列。
然而在文档中它ThreadPoolExecutor说:
当在方法execute(java.lang.Runnable)中提交新任务并且运行的线程数少于corePoolSize时,即使其他工作线程处于空闲状态,也会创建一个新线程来处理该请求。如果运行的线程数超过 corePoolSize 但少于 maxPoolSize ,则仅当队列已满时才会创建新线程。
那么corePoolSize = 0在这种情况下是如何工作的呢?最初,有 0 个线程,因此尽管文档中没有说明,但我认为它将为提交的第一个任务创建一个新线程。但是,现在我们有 1 个线程 > corePoolSize = 0,并且 1 个线程 < MaximumPoolSize = Integer.MAX_VALUE,根据上面的文档“仅当队列已满时才会创建新线程”,但队列是无界的,所以不会再创建新线程,而我们只能使用 1 个线程?
java multithreading executorservice blockingqueue threadpoolexecutor
我正在阅读 LinkedBlockingQueue 代码(JDK8u),我发现 LinkedBlockingQueue 的 head 字段不是私有的,但最后一个字段是私有的。我找不到 head 的任何特定操作。那么为什么不将 head 设置为私有呢?
/**
* Head of linked list.
* Invariant: head.item == null
*/
transient Node<E> head;
/**
* Tail of linked list.
* Invariant: last.next == null
*/
private transient Node<E> last;
Run Code Online (Sandbox Code Playgroud)