【发布时间】:2013-04-27 17:05:20
【问题描述】:
这是my previous question的后续行动:
假设我有一个演员,它每秒处理X 请求。然而,有时会有突发事件,客户端每秒发送Y > X 请求。现在我必须保证客户端在给定的超时时间内收到响应(以太成功或超时状态)。
假设我使用 Scala 和 Akka,你建议如何实现它?
【问题讨论】:
标签: scala concurrency akka actor
这是my previous question的后续行动:
假设我有一个演员,它每秒处理X 请求。然而,有时会有突发事件,客户端每秒发送Y > X 请求。现在我必须保证客户端在给定的超时时间内收到响应(以太成功或超时状态)。
假设我使用 Scala 和 Akka,你建议如何实现它?
【问题讨论】:
标签: scala concurrency akka actor
首先,一个显示超时处理的小代码示例:
import akka.actor._
import akka.util.Timeout
import scala.concurrent.duration._
import akka.pattern._
import scala.util._
import java.util.concurrent.TimeoutException
object TimeoutTest {
def main(args: Array[String]) {
val sys = ActorSystem("test-system")
implicit val timeout = Timeout(2 seconds)
implicit val ec = sys.dispatcher
val ref = sys.actorOf(Props[MyTestActor])
val fut = ref ? "foo"
fut onComplete{
case Success(value) => println("Got success")
case Failure(ex:TimeoutException) => println("Timed out")
case Failure(ex) => println("Got other exception: " + ex.getMessage)
}
}
}
class MyTestActor extends Actor{
def receive = {
case _ =>
Thread.sleep(3000)
sender ! "bar"
}
}
您可以在此示例中看到,我指定了 2 秒的 ask 超时,并且我的演员在响应前 3 秒正在休眠。在这种情况下,我总是会得到 Failure 包装 TimeoutException。现在,超时处理不是 Scala 的 Future 类所固有的,但幸运的是,Akka 为他们的 ask 操作添加了超时支持。在后台,当您执行ask 时,Akka 会创建两个Promises;一个可以由响应消息的参与者完成,另一个由HashedWheelTimer 类中的计时器任务完成。然后,Akka 从这两个Promise 实例中获取Futures,并将它们与Future.firstCompletedOf 合并为一个,因此您从? 调用返回的Future 可以通过接收到的参与者的响应来完成消息或超时,无论先发生什么。
【讨论】:
这取决于你如何使用演员。如果您使用“询问”(如actor ? msg),您会收到一个在指定时间后超时的未来。
见http://doc.akka.io/docs/akka/snapshot/scala/futures.html(与演员一起使用)
您可以向未来添加一个 onFailure 挂钩,以便在未来超时时向客户端发送错误响应。
未来的api:http://www.scala-lang.org/api/current/index.html#scala.concurrent.Future
【讨论】: