【问题标题】:Most correct way work with IO/blocking operations in Akka在 Akka 中使用 IO/阻塞操作的最正确方法
【发布时间】:2016-12-22 19:01:24
【问题描述】:

我有一个关于使用 IO 阻塞操作的正确方法的问题,我读到了正确的模型,它将阻塞操作包装到 Future 并使用这个对象而不是在 actor 内部调用阻塞操作。

但我认为这至少有两个解决方案

case class AccountBalance(id: Int, userId: Int, total: Int)

object AccountBalance {
  def getUserAccountBalance(userId: Int)(
      implicit ec: ExecutionContext): Future[AccountBalance] = {
    Future {
      AccountBalance(1, userId, 1)
    }
  }

  def updateAccountBalance(id: Int, total: Int)(implicit ec:   ExecutionContext): Future[AccountBalance] = {
    Future {
      AccountBalance(id, 1, total)
    }
  }

  // Main logic 
  def getAndInc(userId: Int)(implicit ec: ExecutionContext) = {
    getUserAccountBalance(userId).flatMap { balance =>
    updateAccountBalance(balance.id, balance.total + 1)
  }
 }
}

在第一种方法中,我在 actor 内部使用 AccountBalanece.getAndInc 方法:

class Approach1 extends Actor {
  implicit val executionContext = context.dispatcher

  def receive = {
    case Calculate(userId) =>
      AccountBalance.getAndInc(userId) pipeTo sender
  }
}

另一种解决方案(对我来说更舒服)

class Approach2 extends Actor {
  implicit val executionContext = context.dispatcher

  var firstSender: ActorRef = null
  var userId: Int = -1

  def receive = {
    case Calculate(givenUserId) =>
      userId = givenUserId
      firstSender = sender

      AccountBalance.getUserAccountBalance(userId) pipeTo self

    case AccountBalance(id, _, total) =>
      AccountBalance.updateAccountBalance(id, total + 1) pipeTo    firstSender

  }
} 

哪种解决方案更好(或既危险又不可用)?

关于ExecutionContext 的第二个问题,例如我在实际应用中使用context.dispatcher,我使用下一个:

class A extends Actor {

  implicit val executionContext = ExecutionContext.fromExecutorService(Executors.newFixedThreadPool(10)) 

  override def postStop = {
     executionContext.shutdown()
  }

  def receive = { case _ => }

}

如果我在actor停止后使用executionContext.shutdown,这是关闭线程池并释放所有资源的正确方法吗?

【问题讨论】:

  • 你能清理你的示例代码吗?您的问题是指“阻塞操作”,但 AccountBalance.getUserAccountBalanceAccountBalance.updateAccountBalance 都不是阻塞的,因此根本不需要 Futures...
  • 这只是一个例子,你可以把所有的方法都看作是阻塞的方法

标签: akka


【解决方案1】:

正如反应式编程定义的那样,需要确保执行线程不被阻塞。因此,任何阻塞操作都应该由不同的执行上下文处理,以免阻塞调度参与者命令/消息的线程。我使用第二种模式,我的演员使用“成为”并进入不同的行为,等待来自阻塞操作的响应。这是希望不处理任何新请求并丢失原始发送者actorRef。如果您仍然希望参与者接收其他消息并发送这些消息,您需要在原始发送者和收到的响应之间建立链接。

是的,停止它,这样会起作用,但我建议在创建线程池时要非常小心,并在阻塞参与者之间共享它们。请注意,如果关闭引发异常并且您只想关闭actor A,而不是所有akka系统,则主管不会重新启动actor A。

【讨论】:

    猜你喜欢
    • 2019-10-28
    • 1970-01-01
    • 2019-10-18
    • 2015-08-08
    • 1970-01-01
    • 2019-11-21
    • 1970-01-01
    • 2017-07-16
    • 2012-01-07
    相关资源
    最近更新 更多