【问题标题】:Correct pattern for accumulating state in an Akka actor在 Akka Actor 中累积状态的正确模式
【发布时间】:2014-07-11 23:36:54
【问题描述】:

问题:

在 Akka Actor 中累积状态的正确模式是什么?

上下文:

假设我有一些服务都返回数据。

class ServiceA extends Actor {
  def receive = {
    case _ => sender ! AResponse(100)
  }
}

class ServiceB extends Actor {
   def receive = {
     case _ => sender ! BResponse("n")
   }
}

// ...

我希望有一个控制/监督参与者协调与所有这些服务的对话并跟踪它们的响应,然后将包含所有数据的响应发送回原始发送者。

class Supervisor extends Actor {
  def receive = {
    case "begin" => begin
    case AResponse(id) => ???
    case BResponse(letter) => ???
  }

 // end goal:
 def gotEverything(id: Int, letter: String) =
   originalSender ! (id, letter)

  def begin = {
    ServiceA ! "a"
    ServiceB ! "b"
  }
}

随着服务响应的到来,我如何将所有这些状态关联在一起?据我了解,如果我要将 AResponse 的值分配给 var aResponse: Int,则该 var 会随着接收到不同的消息而不断变化,我无法指望 var 在等待时保持不变对于BResponse 消息。

我意识到我可以使用 ask 和嵌套/flatMap Future,但从我读到的内容来看,这是一个糟糕的模式。有没有办法在没有 Future 的情况下实现这一切?

【问题讨论】:

  • 为什么aResponse 会多次更改?您只向ServiceA 发送一条消息。你的目标到底是什么?等到您收到AResponseBResponse 并用它们的值调用gotEverything
  • >> 你的目标到底是什么?等到您收到 AResponse 和 BResponse 并使用它们的值调用 gotEverything?是的
  • 您可以将响应存储在带有一些标识符或演员参考的列表中

标签: scala akka


【解决方案1】:

因为从不会同时从多个线程访问 Actor,所以您可以轻松地在其中存储和改变您想要的任何状态。例如,您可以这样做:

class Supervisor extends Actor {
  private var originalSender: Option[ActorRef] = None
  private var id: Option[WhateverId] = None
  private var letter: Option[WhateverLetter] = None

  def everythingReceived = id.isDefined && letter.isDefined

  def receive = {
    case "begin" =>
      this.originalSender = Some(sender)
      begin()

    case AResponse(id) =>
      this.id = Some(id)
      if (everythingReceived) gotEverything()

    case BResponse(letter) =>
      this.letter = Some(letter)
      if (everythingReceived) gotEverything()
  }

  // end goal:
  def gotEverything(): Unit = {
    originalSender.foreach(_ ! (id.get, letter.get))
    originalSender = None
    id = None
    letter = None
  }

  def begin(): Unit = {
    ServiceA ! "a"
    ServiceB ! "b"
  }
}

不过,还有更好的方法。您可以使用没有显式状态变量的参与者模拟某种状态机。这是使用become() 机制完成的。

class Supervisor extends Actor {
  def receive = empty

  def empty: Receive = {
    case "begin" =>
      AService ! "a"
      BService ! "b"
      context become noResponses(sender)
  }

  def noResponses(originalSender: ActorRef): Receive = {
    case AResponse(id) => context become receivedId(originalSender, id)
    case BResponse(letter) => context become receivedLetter(originalSender, letter)
  }

  def receivedId(originalSender: ActorRef, id: WhateverId): Receive = {
    case AResponse(id) => context become receivedId(originalSender, id)
    case BResponse(letter) => gotEverything(originalSender, id, letter)
  }

  def receivedLetter(originalSender: ActorRef, letter: WhateverLetter): Receive = {
    case AResponse(id) => gotEverything(originalSender, id, letter)
    case BResponse(letter) => context become receivedLetter(originalSender, letter)
  }

  // end goal:
  def gotEverything(originalSender: ActorRef, id: Int, letter: String): Unit = {
    originalSender ! (id, letter)
    context become empty
  }
}

这可能稍微冗长一些,但它不包含显式变量;所有状态都隐含在Receive方法的参数中,当需要更新这个状态时,actor的receive函数只是切换来反映这个新状态。

请注意,上面的代码非常简单,当有很多“原始发件人”时,它将无法正常工作。在这种情况下,您必须为所有消息添加一个 id 并使用它们来确定哪些响应属于哪个“原始发件人”状态,或者您可以创建多个参与者,每个参与者对应每个“原始发件人”。

【讨论】:

    【解决方案2】:

    我相信 Akka 的方式是使用 actor-per-request 模式。这样,每次收到请求时,您都会创建一个新的参与者,而不是弄清楚哪个响应对应于什么。这是非常便宜的,事实上,每次你做 ask() 时都会发生。

    这些请求处理器(我就是这么称呼它们的)通常具有简单的响应字段。并且只需要简单的 null 比较即可查看请求是否已到达。

    使用此方案,重试/失败也变得更加容易。超时也是如此。

    【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-01-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多