【问题标题】:How do I determine that all actors have received a broadcast message如何确定所有演员都收到了广播消息
【发布时间】:2014-11-23 07:26:04
【问题描述】:

我有一个 ActorA,它从输入流中读取数据并将消息发送到一组 ActorB。当 ActorA 到达输入流的末尾时,它会清理其资源,向 ActorB 广播 Done 消息,然后自行关闭。

我有大约 12 个 ActorB 向一组 ActorC 发送消息。当 ActorB 收到来自 ActorA 的 Done 消息时,它会清理其资源并自行关闭,但最后一个幸存的 ActorB 除外,它在关闭之前向 ActorC 广播 Done 消息。

我有大约 24 个 ActorC 向单个 ActorD 发送消息。与 ActorB 类似,当每个 ActorC 收到 Done 消息时,它会清理其资源并自行关闭,但最后一个幸存的 ActorC 除外,它会向 ActorD 发送 Done 消息。

当 ActorD 收到 Done 消息时,它会清理其资源并自行关闭。

最初我让 ActorB 和 ActorC 在收到 Done 消息后立即传播,但这可能会导致 ActorC 在所有 ActorB 完成处理他们的队列之前关闭;同样,ActorD 可能会在 ActorC 完成队列处理之前关闭。

我的解决方案是使用 ActorB 之间共享的 AtomicInteger

class ActorB(private val actorCRouter: ActorRef,
             private val actorCount: AtomicInteger) extends Actor {
  private val init = {
    actorCount.incrementAndGet()
    ()
  }

  def receive = {
    case Done => {
      if(actorCount.decrementAndGet() == 0) {
        actorCRouter ! Broadcast(Done)
      }
      // clean up resources
      context.stop(self)
    }
  }
}

ActorC 使用类似的代码,每个 ActorC 共享一个 AtomicInteger。

目前所有actor都在一个web service方法中初始化,下游ActorRef在上游actor的构造函数中传入。

有没有首选的方法来做到这一点,例如使用对 Akka 方法的调用而不是 AtomicInteger?


编辑:我正在考虑以下作为可能的替代方案:当演员收到完成消息时,它将接收超时设置为 5 秒(程序将需要一个多小时才能运行,因此延迟清理/关闭几个秒不会影响性能);当 actor 获得 ReceiveTimeout 时,它会向下游 actor 广播 Done、清理并关闭。 (ActorB 和 ActorC 的路由器使用的是 SmallestMailboxRouter)

class ActorB(private val actorCRouter: ActorRef) extends Actor {

  def receive = {
    case Done => {
      context.setReceiveTimeout(Duration.create(5, SECONDS))
    }

    case ReceiveTimeout => {
      actorCRouter ! Broadcast(Done)
      // clean up resources
      context.stop(self)
    }
  }
}

【问题讨论】:

  • 像这样在参与者之间有任何类型的共享状态是一个非常糟糕的主意。将生命周期关注点分开,并确保您的参与者遵守单一职责原则。这样做会容易很多。

标签: algorithm scala concurrency akka actor


【解决方案1】:

在相关参与者之间共享 actorCount 并不是一件好事。 Actor 应该只使用它自己的状态来处理消息。 对于 ActorB 类型的演员,让 ActorBCompletionHanlder 演员怎么样。所有 ActorB 都将引用 ActorBCompletionHanlder 演员。每次 ActorB 收到 Done 消息时,它都可以进行必要的清理并将完成消息传递给 ActorBCompletionHanlder。 ActorBCompletionHanlder 将维护状态变量以维护计数。每次收到完成消息时,它都可以简单地更新计数器。因为这只是这个actor的状态变量,所以不需要它是原子的,这样就不需要任何显式锁定。 ActorBCompletionHanlder 将在收到最后一个完成消息后向 ActorC 发送完成消息。 这种方式的 activeCount 共享不在演员之间,而仅由 ActorBCompletionHanlder 管理。其他类型可以重复相同的操作。

A-> B's -> BCompletionHanlder -> C's -> CCompletionHandler -> D

其他方法可能是为每一组相关的参与者设置一个监控参与者。并且使用监视器上的 watch api 和 child 终止事件,您可以选择在收到 last done 消息后决定要做什么。

val child = context.actorOf(Props[ChildActor])
    context.watch(child)

    case Terminated(child) => {
        log.info(child + " Child actor terminated")
    } 

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-11-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-04-23
    相关资源
    最近更新 更多