【发布时间】: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