【问题标题】:How to Cancel an Akka actor?如何取消 Akka 演员?
【发布时间】:2011-11-17 06:10:48
【问题描述】:

我有一个接收请求并回复它的 akka actor(worker)。请求处理可能需要 3-60 分钟。来电者(也是演员)目前正在使用!!!并等待future.get,但是如果需要,可以更改调用者角色的设计。另外,我目前正在使用 EventDriven 调度程序。

我如何取消(用户发起)请求处理,以便释放工作角色并返回就绪状态以接收新请求?我希望有一个类似于 java.util.concurrent.Future 的取消方法的方法,但在 Akka 1.1.3 中找不到

编辑:

我们试图通过completeWithException 获得我们正在寻找的行为:

object Cancel {
  def main(args: Array[String]) {
    val actor = Actor.actorOf[CancelActor].start
    EventHandler.info(this, "Getting future")
    val future = (actor ? "request").onComplete(x => EventHandler.info(this, "Completed!! " + x.get))
    Thread.sleep(500L)
    EventHandler.info(this, "Cancelling")
    future.completeWithException(new Exception("cancel"))
    EventHandler.info(this, "Future is " + future.get)
  }
}

class CancelActor extends Actor {
  def receive = {
    case "request" =>
      EventHandler.info(this, "start")
      (1 to 5).foreach(x => {
        EventHandler.info(this, "I am a long running process")
        Thread.sleep(200L)
      })
      self reply "response"
      EventHandler.info(this, "stop")
  }
}

但这并没有停止长期运行的过程。

    [INFO]    [9/16/11 1:46 PM] [main] [Cancel$] Getting future
    [INFO]    [9/16/11 1:46 PM] [akka:event-driven:dispatcher:global-2] [CancelActor] start
    [INFO]    [9/16/11 1:46 PM] [akka:event-driven:dispatcher:global-2] [CancelActor] I am a long running process
    [INFO]    [9/16/11 1:46 PM] [akka:event-driven:dispatcher:global-2] [CancelActor] I am a long running process
    [INFO]    [9/16/11 1:46 PM] [akka:event-driven:dispatcher:global-2] [CancelActor] I am a long running process
    [INFO]    [9/16/11 1:46 PM] [main] [Cancel$] Cancelling
    [ERROR]   [9/16/11 1:46 PM] [akka:event-driven:dispatcher:global-7] [ActorCompletableFuture] 
    java.lang.Exception: cancel
        at kozo.experimental.Cancel$.main(Cancel.scala:15)
...

    [INFO]    [9/16/11 1:46 PM] [akka:event-driven:dispatcher:global-2] [CancelActor] I am a long running process
    [INFO]    [9/16/11 1:46 PM] [akka:event-driven:dispatcher:global-2] [CancelActor] I am a long running process
    [INFO]    [9/16/11 1:46 PM] [akka:event-driven:dispatcher:global-2] [CancelActor] stop

相比之下,考虑 java.util.concurrent.Future 的行为:

object Cancel2 {
  def main(args: Array[String]) {
    val executor: ExecutorService = Executors.newSingleThreadExecutor()
    EventHandler.info(this, "Getting future")
    val future = executor.submit(new Runnable {
      def run() {
        EventHandler.info(this, "start")
        (1 to 5).foreach(x => {
          EventHandler.info(this, "I am a long running process")
          Thread.sleep(200L)
        })
      }
    })
    Thread.sleep(500L)
    EventHandler.info(this, "Cancelling")
    future.cancel(true)
    EventHandler.info(this, "Future is " + future.get)
  }
}

哪个会停止长时间运行的进程

    [INFO]    [9/16/11 1:48 PM] [main] [Cancel2$] Getting future
    [INFO]    [9/16/11 1:48 PM] [pool-1-thread-1] [anon$1] start
    [INFO]    [9/16/11 1:48 PM] [pool-1-thread-1] [anon$1] I am a long running process
    [INFO]    [9/16/11 1:48 PM] [pool-1-thread-1] [anon$1] I am a long running process
    [INFO]    [9/16/11 1:48 PM] [pool-1-thread-1] [anon$1] I am a long running process
    Exception in thread "main" java.util.concurrent.CancellationException
...
    [INFO]    [9/16/11 1:48 PM] [main] [Cancel2$] Cancelling

【问题讨论】:

  • 不清楚是否是同一段代码提交了希望能够取消它的作业。我认为如果您能描述实际的业务问题会有所帮助。
  • 调用者在 future.get 上被阻止,所以我无法用它做任何事情。在这一点上问题是抽象的:worker actor 正在处理一个请求,而取消/中断 actor 的优雅方式将是什么。由来电者/其他演员...

标签: akka actor


【解决方案1】:

您还可以在 Actor 中检查 Future 的状态。

class MyActor extends Actor {
  def receive = {
    case msg =>
      while(!self.senderFuture.get.isCompleted) {
        performWork(msg)
      }
      self reply result
  }
  ...
}

这需要使用“?”发送消息或“问”。 希望能帮助到你。

【讨论】:

    【解决方案2】:

    如果您只是在虚拟机中,您可以将 AtomicBoolean 与您的 Job 消息一起传递,并间歇性地检查您的 actor 以查看您是否应该中止。

    actor ! Job(..., someAtomicBoolean)
    
    class MyActor extends Actor {
      def receive = {
        case Job(..., cancelPlease) =>
          while(cancelPlease.get == false) {
            performWork
          }
          self reply result
      }
    }
    

    【讨论】:

    • NOOOOOOOOOOOo00oooo0o0o0!! “消息可以是任何类型的对象,但必须是不可变的”doc.akka.io/docs/akka/snapshot/java/untyped-actors.html 此消息不会是不可变的!
    • 你不是也在开发akka吗?很惊讶看到你提出这样的建议。请解释一下?
    猜你喜欢
    • 1970-01-01
    • 2018-09-29
    • 2014-12-02
    • 2019-04-23
    • 2013-11-22
    • 1970-01-01
    • 1970-01-01
    • 2013-11-29
    • 2012-10-09
    相关资源
    最近更新 更多