我遇到了这个scala.concurrent.blocking方法,根据Scala文档,这是......
用于指定可能阻塞的代码段,允许当前的BlockContext调整运行时的行为.正确标记阻塞代码可以提高性能或避免死锁.
我有些疑惑:
scala.concurrent.ExecutionContext.Implicits.global执行上下文或用户创建的执行上下文?blocking {... 包装任何可执行文件会发生什么}?Future(blocking(blockingCall()))和之间有什么区别blocking(Future(blockingCall()))?这两个都在中定义scala.concurrent._
我使用spray-can,spray-http 1.3.2和akka 2.3.6设置了一个http服务器.我的application.conf没有任何akka(或spray)条目.我的演员代码:
class TestActor extends HttpServiceActor with ActorLogging with PlayJsonSupport {
val route = get {
path("clientapi"/"orders") {
complete {{
log.info("handling request")
System.err.println("sleeping "+Thread.currentThread().getName)
Thread.sleep(1000)
System.err.println("woke up "+Thread.currentThread().getName)
Seq[Int]()
}}
}
}
override def receive: Receive = runRoute(route)
}
Run Code Online (Sandbox Code Playgroud)
像这样开始:
val restService = system.actorOf(Props(classOf[TestActor]), "rest-clientapi")
IO(Http) ! Http.Bind(restService, serviceHost, servicePort)
Run Code Online (Sandbox Code Playgroud)
当我发送10个并发请求时,它们都被喷射立即接受并转发给不同的调度程序参与者(根据我从applicaiton.conf中删除的akka的日志配置,以免影响结果),但所有都由同一个线程处理,睡觉,只有在醒来之后才会收到下一个请求.
我应该在配置中添加/更改什么?从我在reference.conf中看到的,默认执行程序是fork-join-executor,所以我希望所有请求都是开箱即用的并行执行.
我们正在为每个(小)传入的消息组创建一个actor链,以保证它们的顺序处理和管道(组通过公共id区分).问题是,我们的链叉,就像A1 -> (A2 -> A3 | A4 -> A5)我们要保证消息经历之间没有种族A2 -> A3和A4 -> A5.传统的解决方案是阻止A1actor直接当前消息被完全处理(在一个子链中):
def receive { //pseudocode
case x => ...
val f = A2orA4 ? msg
Await.complete(f, timeout)
}
Run Code Online (Sandbox Code Playgroud)
因此,应用程序中的线程数与处理中的消息数成正比,无论这些消息是活动的还是异步等待来自外部服务的某些响应.它使用fork-join(或任何其他动态)池工作大约两年,但当然不能使用固定池,并且在高负载的情况下极大地降低性能.更重要的是,它会影响GC,因为每个被阻塞的fork-actor都会保留冗余的先前消息的状态.
即使使用背压,它也会创建比收到的消息多N倍的线程(因为流中有N个顺序叉),这仍然很糟糕,因为一条消息的处理需要很长时间但CPU不多.所以我们应该处理更多的消息,因为我们有足够的内存.我提出的第一个解决方案 - 将链条线性化A1 -> A2 -> A3 -> A4 -> A5.还有更好的吗?