【发布时间】: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