【问题标题】:Initializing an akka actor using Akka-IO TCP使用 Akka-IO TCP 初始化 akka actor
【发布时间】:2014-09-29 14:59:52
【问题描述】:

使用Akka-IO TCP,在actor中建立连接的过程如下:

class MyActor(remote: InetSocketAddress) extends Actor {

  IO(Tcp) ! Connect(remote)    //this is the first step, remote is the address to connect to

  def receive = {
    case CommandFailed(_: Connect) => context stop self // failed to connect

    case Connected(remote, local) =>
      val connection = sender()
      connection ! Register(self)
      // do cool things...  
  }
}

您向IO(Tcp) 发送Connect 消息并期望收到CommandFailed 或Connected 消息。

现在,我的目标是创建一个包装 TCP 连接的参与者,但我希望我的参与者仅在建立连接后才开始接受消息 - 否则,在等待 Connected 消息时,它将开始接受查询,但没有人可以寄给他们。

我尝试了什么:

class MyActor(address: InetSocketAddress) extends Actor {

  def receive = {
    case Initialize =>
      IO(Tcp) ! Connect(address)
      context.become(waitForConnection(sender()))

    case other => sender ! Status.Failure(new Exception(s"Connection to $address not established yet."))
  }

  private def waitForConnection(initializer: ActorRef): Receive = {
    case Connected(_, _) =>
      val connection = sender()
      connection ! Register(self)
      initializer ! Status.Success(Unit)
      // do cool things

    case CommandFailed(_: Connect) =>
      initializer ! Status.Failure(new Exception("Failed to connect to " + host))
      context stop self
  }
}

我的第一个receive 期待一个虚构的Initialize 消息,它将触发整个连接过程,一旦完成,Initialize 的发送者会收到一条成功消息并知道它可以知道开始发送查询。

我对此不太满意,它迫使我用它来创建我的演员

val actor = system.actorOf(MyActor.props(remote))
Await.ready(actor ? Initialize, timeout)

而且它不会很“重启”友好。

在 Tcp 层回复 Connected 之前,有什么办法可以保证我的 actor 不会开始从邮箱接收消息吗?

【问题讨论】:

    标签: scala akka akka-io


    【解决方案1】:

    使用Stash trait 来存储您现在无法处理的消息。当每个过早的消息到达时,使用stash() 来推迟它。连接打开后,使用unstashAll() 将这些消息返回到邮箱进行处理。然后可以使用become() 切换到消息处理状态。

    【讨论】:

    • 听起来不错。我正在考虑自己实现这个“缓冲”,不知道有一个内置的特性。谢谢,为我节省了一些工作:)
    【解决方案2】:

    为了让你的actor对重启更友好,你可以重写与actor生命周期相关的方法,例如preStart或postStop。 Akka documentation 对actor的启动、停止和重启钩子有很好的解释。

    class MyActor(remote: InetSocketAddress) extends Actor {
    
      override def preStart() {
        IO(Tcp) ! Connect(remote) 
      }
    
      ...
    }
    

    现在您可以使用val actor = system.actorOf(MyActor.props(remote)) 开始您的演员。它在启动时建立连接,并在重新启动时重新建立。

    【讨论】:

    • 是的,但是在使用actor之前没有办法等待连接完全建立。
    • 你可以同时使用这个和 Bob 的解决方案。
    猜你喜欢
    • 2020-06-29
    • 2023-04-02
    • 2015-06-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-12-19
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多