【问题标题】:Akka pattern for handling asynchronous actions in receive用于在接收中处理异步操作的 Akka 模式
【发布时间】:2014-05-13 06:00:15
【问题描述】:

我有一个 Actor,它接收指标数据点并定期聚合并将它们保存到磁盘。后一个操作执行 I/O,所以我不想使用阻塞操作。但是如果我将它切换为异步,我如何防止在聚合完成之前接收到其他数据点而不阻塞某处。

我见过的一种模式是使用Stash,如下所示:

class Aggregator extends Actor with Stash {
  def receive = processing

  def processing: Receive = {
    case "aggregate" => {
      context.become(aggregating)
      aggregate().onComplete {
        case Success => self ! "aggregated"
        case Failure => self ! "aggregated"
      }
    }
    case msg => ??? // Process task
  }

  def aggregating: Receive = {
    case "aggregated" =>
      unstashAll()
      context.become(processing)
    case msg =>
      stash()
  }
}

我对此的疑虑是,我的聚合操作的完成只是任何人都可以发送的消息。据我了解,我无法在未来的完成中影响“不合时宜”。

作为旁注,我无法确定像 onComplete 这样的完成是否由与 receive 相同的调度程序以某种方式执行,因为如果不是,完成将破坏参与者的单线程保护,否则报价。

或者有没有更好的模式来完成receive 内部不同步和直接的操作,同时保证我的状态在我完成之前不能改变?每当actor状态处理任何类型的I/O(如DB)时,这种情况似乎都是如此,显然您希望尽可能避免同步I/O。

【问题讨论】:

  • 将这些消息转发给一个特殊的演员,他的工作只是保存指标。这样原actor就不会节流,指标会一一保存
  • 解决您对任何人都可以发送消息的担忧。您可以在触发聚合的伴随对象中创建一个私有案例对象。无论如何,其他人将其拆分为多个演员的想法确实值得一试。

标签: scala asynchronous io akka


【解决方案1】:

您的聚合器参与者当前正在做两件事:聚合和存储。您可以通过拆分这两个任务来解决您的问题并简化您的系统。 single-responsibility-principle 也适用于演员。

我会创建一个专门的 actor 用于编写,并创建一个消息类来保存聚合数据。这个actor子系统应该是这样的:

理想情况下,写入磁盘所需的时间比聚合间隔要短,这样您的系统才能保持稳定。在出现峰值的情况下,DataStore Actor 的队列将作为要写入存储的消息的缓冲区。

根据您的应用程序,您可能需要实现某种形式的 ack & retries 以确保已写入聚合数据。

【讨论】:

  • 你是对的,我唯一关心防止接收新指标的时候是在我收集要聚合的值时,但聚合和持久性不需要阻止接收指标,所以聚合消息是一种快速的内存操作,将 的值发送给另一个参与者。仍在学习如何分解我常用的对象模型,并将依赖项注入到 actor 中。
猜你喜欢
  • 2019-11-30
  • 1970-01-01
  • 1970-01-01
  • 2019-10-16
  • 2011-10-03
  • 1970-01-01
  • 1970-01-01
  • 2015-11-22
  • 1970-01-01
相关资源
最近更新 更多