【问题标题】:Akka Actor running infinite loop after initialization messageAkka Actor 在初始化消息后运行无限循环
【发布时间】:2016-04-07 14:32:43
【问题描述】:

Akka 和 Actors 的新手 - 我需要启动一些 Actor,这些 Actor 基本上会花费一生的时间阅读 Kafka 主题并写入 Ignite 缓存。我这样配置调度程序:

kafka-dispatcher {
  executor  = "thread-pool-executor"
  type      = PinnedDispatcher
}

我的演员是用.withDispatcher("kafka-dispatcher") 创建的,我的假设是每个演员都会被分配一个单独的线程。

这些演员基本上都是这样度过一生的:

override def receive: Receive = LoggingReceive {
  case InitWorker => {
    initialize()
    pollTopic() // This never returns
  }
}

换句话说,它们接收到初始化消息,然后调用pollTopic() 方法,该方法永远不会返回 - 它运行循环读取(在有数据之前将阻塞)然后写入数据。

我的问题:

  1. 这是犹太洁食吗?
  2. 是否有更好的方式,即更惯用的方式来做到这一点?请注意,pollTopic() 内的读取调用会阻塞。

【问题讨论】:

    标签: scala akka actor


    【解决方案1】:

    回答您的第 2 点并根据您想要做的事情的描述,也许您想考虑将 Akka streamsreactive-kafka 库一起使用。 Akka 流在底层使用 Actor,但会为您管理所有这些,因此您可以只专注于实现只做一件事的可重用小组件。

    然后您将能够编写数据处理管道,使用 Kafka 作为您的数据流的Source。我对 Ignite 缓存了解不多,但您可能会为它写一个 Sink 或者 - 如果您正在谈论阻塞 API,mapAsync 将是您的朋友。

    【讨论】:

      猜你喜欢
      • 2014-09-29
      • 1970-01-01
      • 2020-06-29
      • 2017-04-21
      • 1970-01-01
      • 2016-04-17
      • 1970-01-01
      • 1970-01-01
      • 2013-04-30
      相关资源
      最近更新 更多