【问题标题】:Akka: Guarding an async resourceAkka:保护异步资源
【发布时间】:2014-05-14 17:08:02
【问题描述】:

虽然我的 SRP 方式的错误已得到纠正 yesterday,但我仍然想知道您如何干净地保证单线程访问 akka 中的异步资源,例如文件句柄。显然我不想允许从不同的线程对它分派多个读取和写入操作,但是如果我的参与者在该文件上调用基于未来的 API,这很可能会发生。

我想出的最好的模式是这样的:

trait AsyncIO {
  def Read(offset: Int, count: Int) : Future[ByteBuffer] = ???
}

object GuardedIOActor {
  case class Read(offset: Int, count: Int)
  case class ReadResult(data: ByteBuffer)
  private case class ReadCompleted()
}

class GuardedIOActor extends Actor with Stash with AsyncIO {
  import GuardedIOActor._
  var caller :Option[ActorRef] = None

  def receive = {
    case Read(offset,count) => caller match {
      case None => {
        caller = Some(sender)
        Read(offset,count).onSuccess({
          case data => {
            self ! ReadCompleted()
            caller.get ! ReadResult(data)
          }
        })
      }
      case Some(_) => {
        stash()
      }
    }
    case ReadCompleted() => {
      caller = None
      unstashAll()
    }
  }
}

但是这个要求对我来说还不够深奥。我的意思是应该有大量的资源需要同步访问但有一个异步 API。我是否忽略了一些常见的命名模式?

【问题讨论】:

  • 也许this 答案和akka 文档中描述的pipeTo 方法可以帮助您。
  • 如果我理解期货的调度和正确的管道,我认为这对我没有帮助。两个 Read 消息一个接一个地进入,都创建期货,它们通过 ExecutionContext 分派,并且可以非常可行地在两个不同的线程上执行,都试图同时访问同一个文件句柄。我的意思是,如果 AsynchronousFileChannel 已经保证同步访问,这一切都没有实际意义,但我似乎找不到这样的保证。即使是这样,这也只有在我的操作只需要在频道上进行一次调用时才成立
  • 我明白了,我在想你是否真的愿意将文件操作交给多线程。我认为你应该有一个 Actor 来序列化对文件的访问,这意味着只有一个 Actor(并且以这种方式一个线程)可以访问资源,然后在接收和执行操作时有一个邮箱,可能有一个将异步 IO 操作包装在可以多次打开的未来中通常是一个糟糕的设计理念,正如您所说,这会破坏资源访问管理逻辑。
  • 这是一个序列化对文件的所有访问的actor。我正在使用未来,因为参与者在执行 IO 时不应该阻塞线程,并且通过使用异步 IO API,我放弃了对完成返回的线程的控制,作为不阻塞线程的奖励演员系统
  • 当然可以,但是您有机会打开多个异步 IO 操作,我的意思是 Actor 应该序列化访问和操作确保只有一个一个线程打开一个文件。丹的答案正是我的想法,如果我不够清楚,请见谅。

标签: scala asynchronous io akka future


【解决方案1】:

我认为您的解决方案的要点还不错,但是您可以使用 context.become 使您的演员表现得更像一个状态机:

class GaurdedIOActor extends Actor with Stash with AsyncIO {
  import GuardedIOActor._

  def receive = notReading

  def notReading: Receive = {
    case Read(offset, count) => {
      val caller = sender
      Read(offset,count).onSuccess({
        case data => {
          self ! ReadCompleted()
          caller ! ReadResult(data)
        }
      })
      context.become(reading)
    }
  }

  def reading: Receive = {
    case r: Read => stash()
    case ReadCompleted() => {
      context.become(notReading)
      unstashAll()
    }
  }
}

现在你的演员有两个明确定义的状态,不需要var

【讨论】:

  • 我避免成为/不成为,因为它添加了更多样板,但是看看这两个实现,你的例子更容易阅读和推理
  • 所以我想这是最好的模式。我仍然感到惊讶,因为对于我认为相当基本的用例来说,这是很多手动参与者管理
【解决方案2】:

我意识到这个添加已经晚了一年,但是这个问题帮助我推理了我遇到的类似情况。以下 trait 封装了上面建议的功能,以减少参与者中的样板。它可以与任何带有 Stash 的演员混合。用法类似于pipeTo 模式;只需输入future.pipeSequentiallyTo(sender),您的actor 将在future 完成并发送响应之前不会处理消息。

import scala.concurrent.ExecutionContext
import scala.concurrent.Future
import scala.language.implicitConversions
import scala.util.Failure
import scala.util.Success

import akka.actor.Actor
import akka.actor.ActorRef
import akka.actor.Stash
import akka.actor.Status

trait SequentialPipeToSupport { this: Actor with Stash =>

  case object ProcessingFinished

  def gotoProcessingState() = context.become(processingState)

  def markProcessingFinished() = self ! ProcessingFinished

  def processingState: Receive = {
    case ProcessingFinished =>
      unstashAll()
      context.unbecome()
    case _ =>
      stash()
  }

  final class SequentialPipeableFuture[T](val future: Future[T])(implicit executionContext: ExecutionContext) {

    def pipeSequentiallyTo(recipient: ActorRef): Future[T] = {

      gotoProcessingState()

      future onComplete {
        case Success(r) =>
          markProcessingFinished()
          recipient ! r
        case Failure(f) =>
          markProcessingFinished()
          recipient ! Status.Failure(f)
      }

      future
    }
  }

  implicit def pipeSequentially[T](future: Future[T])(implicit executionContext: ExecutionContext) = 
    new SequentialPipeableFuture(future)

}

也可作为Gist使用

【讨论】:

  • 看起来很有趣。我需要玩它
猜你喜欢
  • 2021-02-15
  • 1970-01-01
  • 2015-08-10
  • 2017-08-18
  • 2015-09-05
  • 2014-01-27
  • 2017-01-17
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多