【问题标题】:Can I safely create a Thread in an Akka Actor?我可以在 Akka Actor 中安全地创建线程吗?
【发布时间】:2016-07-20 16:37:08
【问题描述】:

我有一个 Akka Actor,我想向其发送“控制”消息。 这个 Actor 的核心任务是监听 Kafka 队列,这是一个循环内的轮询过程。

我发现以下只是简单地锁定了 Actor 并且它不会收到“停止”(或任何其他)消息:

class Worker() extends Actor {
  private var done = false

  def receive = {
    case "stop" => 
      done = true
      kafkaConsumer.close()
    // other messages here
  }

  // Start digesting messages!
  while (!done) {
    kafkaConsumer.poll(100).iterator.map { cr: ConsumerRecord[Array[Byte], String] =>
      // process the record
      ), null)
    }
  }
}

我可以将循环包装在由 Actor 启动的线程中,但是从 Actor 内部启动线程是否可以/安全?有没有更好的办法?

【问题讨论】:

  • 绝对不是!您可能想要创建类似“消费者演员”的东西。看看Reactive Kafka

标签: scala akka


【解决方案1】:

基本上你可以,但请记住,这个演员会被阻止,并且规则的拇指是永远不要阻止内部演员。如果您仍想这样做,请确保此 Actor 在与本机线程池不同的线程池中运行,这样您就不会影响 Actor 系统的性能。另一种方法是向自身发送消息以轮询新消息。

1) 接收来自 kafka 的轮询消息的命令

2) 移交 给相关参与者的消息

3) 向自己发送消息以订购 拉新消息

4) 交出...

代码明智:

case object PollMessage

class Worker() extends Actor {
  private var done = false

  def receive = {
    case PollMessage ⇒ {
      poll()
      self ! PollMessage
    }
    case "stop" =>
      done = true
      kafkaConsumer.close()
    // other messages here
  }

  // Start digesting messages!

  def poll() = {
    kafkaConsumer.poll(100).iterator.map { cr: ConsumerRecord[Array[Byte], String] =>
      // process the record
      ), null)
    }
  }

}

我不确定如果你持续阻止演员,你是否会收到停止消息。

【讨论】:

    【解决方案2】:

    添加@Louis F. 答案;根据您的演员的配置,如果在给定的时刻他们很忙,他们将丢弃他们收到的所有消息,或者将它们放入邮箱,即队列,稍后将处理这些消息(通常以先进先出的方式)。但是,在这种特殊情况下,您将使用PollMessage 淹没演员,并且您无法保证您的消息不会被丢弃 - 这似乎发生在您的情况中。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2010-10-26
      • 2010-11-29
      • 1970-01-01
      • 2022-01-10
      • 2018-07-17
      • 1970-01-01
      相关资源
      最近更新 更多