【发布时间】:2015-07-17 02:54:21
【问题描述】:
更新:我发现如果我将ThreadPoolExecutor's 核心池大小设置为与最大池大小(29 个线程)相同,我的程序仍然可以响应。但是,如果我将核心池大小设置为 11,将最大池大小设置为 29,那么参与者系统只会创建 11 个线程。如何配置 ActorSystem / ThreadPoolExecutor 以继续创建线程以超过核心线程数并保持在最大线程数内?我不希望将核心线程数设置为最大线程数,因为我只需要额外的线程来取消作业(这应该是一种罕见的事件)。
我有一个针对 Oracle 数据库运行的批处理程序,使用 Java/Akka 类型的 Actor 和以下 Actor 实现:
-
BatchManager负责与 REST 控制器通信。它管理Queue的未初始化批处理作业;当从队列中轮询一个未初始化的批处理作业时,它会变成一个JobManageractor 并执行。 -
JobManager维护一个存储过程队列和一个Workers池;它使用存储过程初始化每个Worker,当Worker完成时,它会将过程的结果发送到JobManager,然后JobManager将另一个存储过程发送到Worker。当作业队列为空且所有Workers都空闲时,批处理终止,此时JobManager将其结果报告给BatchManager,关闭其工作人员(通过TypedActor.context().stop()),然后自行关闭。JobManager有一个Promise<Status> completion,当作业成功完成或作业因取消或致命异常终止时完成。 -
Worker执行存储过程。它创建OracleConnection 和一个CallableStatement 用于执行存储过程,并将onFailure回调与JobManager.completion注册到abort连接和cancel语句。此回调不使用 Actor 系统的执行上下文,而是使用从BatchManager中创建的缓存执行器服务创建的执行上下文。
我的配置是
{"akka" : { "actor" : { "default-dispatcher" : {
"type" : "Dispatcher",
"executor" : "default-executor",
"throughput" : "1",
"default-executor" : { "fallback" : "thread-pool-executor" }
"thread-pool-executor" : {
"keep-alive-time" : "60s",
"core-pool-size-min" : coreActorCount,
"core-pool-size-max" : coreActorCount,
"max-pool-size-min" : maxActorCount,
"max-pool-size-max" : maxActorCount,
"task-queue-size" : "-1",
"task-queue-type" : "linked",
"allow-core-timeout" : "on"
}}}}}
worker数在别处配置,目前workerCount = 8; coreActorCount 是 workerCount + 3 而 maxActorCount 是 workerCount * 3 + 5。我在具有两个内核和 8GB 内存的 Macbook Pro 10 上进行测试;生产服务器要大得多。我正在与之交谈的数据库位于极慢的 VPN 后面。我正在使用 Oracle 的 JavaSE 1.8 JVM 运行所有这些。本地服务器是 Tomcat 7。Oracle JDBC 驱动程序是 10.2 版(我可能能够说服使用更新版本的权力)。所有方法要么返回void 或Future<>,并且应该是非阻塞的。
当一个批次成功终止时,就没有问题了 - 下一个批次立即开始,有完整的工人。但是,如果我通过JobManager#completion.tryFailure(new CancellationException("Batch cancelled")) 终止当前批处理,则Workers 注册的onFailure 回调触发,然后系统变得无响应。调试 printlns 表明新批次以八分之三的正常工作人员开始,BatchManager 变得完全没有响应(我添加了一个 Future<String> ping 命令,它只返回一个Futures.successful("ping"),这也超时了)。 onFailure 回调在单独的线程池中执行,即使它们在参与者系统的线程池中,我也应该有足够高的 max-pool-size 以容纳原始 JobManager、Workers、onFailure回调,第二个JobManager 是Workers。相反,我似乎正在容纳原来的JobManager 和它的Workers,新的JobManager 和不到一半的Workers,而BatchManager. 没有任何剩余我正在运行它的计算机是资源不足,但看起来应该可以运行十几个线程。
这是配置问题吗?这是由于 JVM 强加的限制和/或 Tomcat 强加的限制吗?这是因为我处理阻塞 IO 的方式有问题吗?可能还有其他几件事我可能做错了,这些就是我想到的。
Gist of CancellableStatement CallableStatement 和 OracleConnection 被取消
Gist of Immutable 创建CancellableStatements 的位置
Gist of JobManager's cleanup code
Config dump通过System.out.println(mergedConfig.toString());获得
编辑:我相信我已经将问题缩小到演员系统(无论是它的配置还是它与阻塞数据库调用的交互)。我消除了Worker 演员并将他们的工作负载转移到Runnables,在固定大小的ThreadPoolExecutor 上执行,其中每个JobManager 创建自己的ThreadPoolExecutor 并在批处理完成时将其关闭(shutDown on正常终止,shutDownNow 异常终止)。取消在BatchManager 中实例化的缓存线程池上运行。 Actor 系统的调度程序仍然是ThreadPoolExecutor,但只分配了六个线程。使用这种替代设置,取消按预期执行 - 工作人员在其数据库连接被中止时终止,并且新的JobManager 立即执行完整的工作线程。这向我表明这不是硬件/JVM/Tomcat 问题。
更新:我使用Eclipse's Memory Analyzer 进行了线程转储。我发现取消线程挂在CallableStatement.close() 上,所以我对取消重新排序,使OracleConnection.abort() 在CallableStatement.cancel() 之前,这解决了问题——所有取消(显然)都正确执行。不过,Worker 线程继续挂在他们的声明上 - 我怀疑我的 VPN 可能部分或全部归咎于此。
PerformanceAsync-akka.actor.default-dispatcher-19
at java.net.SocketInputStream.socketRead0(Ljava/io/FileDescriptor;[BIII)I (Native Method)
at java.net.SocketInputStream.read([BIII)I (SocketInputStream.java:150)
at java.net.SocketInputStream.read([BII)I (SocketInputStream.java:121)
at oracle.net.ns.Packet.receive()V (Unknown Source)
at oracle.net.ns.DataPacket.receive()V (Unknown Source)
at oracle.net.ns.NetInputStream.getNextPacket()V (Unknown Source)
at oracle.net.ns.NetInputStream.read([BII)I (Unknown Source)
at oracle.net.ns.NetInputStream.read([B)I (Unknown Source)
at oracle.net.ns.NetInputStream.read()I (Unknown Source)
at oracle.jdbc.driver.T4CMAREngine.unmarshalUB1()S (T4CMAREngine.java:1109)
at oracle.jdbc.driver.T4CMAREngine.unmarshalSB1()B (T4CMAREngine.java:1080)
at oracle.jdbc.driver.T4C8Oall.receive()V (T4C8Oall.java:485)
at oracle.jdbc.driver.T4CCallableStatement.doOall8(ZZZZ)V (T4CCallableStatement.java:218)
at oracle.jdbc.driver.T4CCallableStatement.executeForRows(Z)V (T4CCallableStatement.java:971)
at oracle.jdbc.driver.OracleStatement.doExecuteWithTimeout()V (OracleStatement.java:1192)
at oracle.jdbc.driver.OraclePreparedStatement.executeInternal()I (OraclePreparedStatement.java:3415)
at oracle.jdbc.driver.OraclePreparedStatement.execute()Z (OraclePreparedStatement.java:3521)
at oracle.jdbc.driver.OracleCallableStatement.execute()Z (OracleCallableStatement.java:4612)
at com.util.CPProcExecutor.execute(Loracle/jdbc/OracleConnection;Ljava/sql/CallableStatement;Lcom/controller/BaseJobRequest;)V (CPProcExecutor.java:57)
但是,即使在修复了取消订单之后,我仍然遇到 Actor 系统没有创建足够线程的问题:在新批次中我仍然只获得了八分之三的工作线程,新的工作线程被添加为被取消的工作人员的网络连接超时。我总共有 11 个线程 - 我的核心池大小,29 个线程中 - 我的最大池大小。显然,演员系统忽略了我的最大池大小参数,或者我没有正确配置最大池大小。
【问题讨论】:
-
您是每次在工作线程中创建新连接还是使用 ConnectionPool 中的 getConnection()?如果您在 Worker 线程中创建新连接,请确保该连接已正确释放。如果您即使使用工作线程中的 ConnectionPool 也遇到问题,我们必须以不同的方式解决问题
-
@sunrise76 我每次都使用 try-with-resources 块创建一个新连接,以确保它被关闭
-
创建和关闭连接非常昂贵。请摆脱它并使用带有 getConnection() 和 releaseConnection() 的连接池。 releaseConnection 将在没有 close() 的情况下释放与池的连接,此模型将针对高流量进行扩展。还有一件事 - 获取应用程序的线程转储并注意“等待监视器进入”和“RUNNABLE”和“死锁”语句
-
@sunrise76 Spring 负责在后台缓存连接,因此
close要么将连接返回到池中,要么关闭连接,具体取决于 Spring 的(即核心连接池是否已超过)。我会四处寻找,看看我是否可以转储 Actor 系统的线程/Actor 状态 - 它从所有事物中代理了 bejeesus,因此很难获得它的原始 ExecutorService。 -
你从转储中得到任何线索吗?您的服务器是否在 Linux 中运行?
标签: java multithreading scala akka blocking