【问题标题】:Akka's actor based custom Event Bus implementation causes bottleneckAkka 的基于 actor 的自定义事件总线实现导致瓶颈
【发布时间】:2016-01-16 15:33:41
【问题描述】:

我正在尝试在 Akka 的演员模型之上实现事件总线 (Pub-Sub) 模式。

“Native”EventBus 实现不符合我的一些要求(例如,可能只保留主题中的最后一条消息,它特定于 MQTT 协议,我正在为它实现消息代理https://github.com/butaji/JetMQ)。

我的EventBus当前接口如下:

object Bus {
  case class Subscribe(topic: String, actor: ActorRef)
  case class Unsubscribe(topic: String, actor: ActorRef)
  case class Publish(topic: String, payload: Any, retain: Boolean = false)
}

用法如下:

val system = ActorSystem("System")
val bus = system.actorOf(Props[MqttEventBus], name = "bus")
val device1 = system.actorOf(Props(new DeviceActor(bus)))
val device2 = system.actorOf(Props(new DeviceActor(bus)))

所有的设备都引用了一个总线actor。 Bus Actor 负责存储订阅和主题的所有状态(例如保留消息)。

设备参与者可以决定他们想要发布、订阅或取消订阅主题的任何内容。

经过一些性能基准测试后,我意识到我当前的设计会影响发布和订阅之间的处理时间,原因如下:

  1. 我的 EventBus 实际上是一个单例
  2. 它造成了巨大的处理负载队列

如何为我的事件总线实施分配(并行化)工作负载? 当前的解决方案是否适合 akka-cluster?

目前,我正在通过 Bus 的几个实例考虑routing,如下所示:

val paths = (1 to 5).map(x => {
  system.actorOf(Props[EventBusActor], name = "event-bus-" + x).path.toString
})

val bus_publisher = system.actorOf(RoundRobinGroup(paths).props())
val bus_manager = system.actorOf(BroadcastGroup(paths).props())

地点:

  • bus_publisher 将负责获取 Publish,
  • bus_manager 将负责获取订阅/取消订阅。

如下所示,它将在所有总线上复制状态,并随着负载的分配减少每个参与者的队列。

【问题讨论】:

  • 你有多少订阅者和发布者?
  • @RomainHippeau 每个节点最多几千个,deviceActor 实际上是 TCP 连接的影子。每个 deviceActor 通常可以有几个不同主题的订阅
  • 也许您只是在违反 pub/sub 范式的限制。它的可扩展性不是很好。

标签: scala akka publish-subscribe mqtt akka-cluster


【解决方案1】:

您可以在单件巴士内部而不是外部路线。您的总线可能负责路由消息和建立主题,而子 Actor 可能负责分发消息。一个基本示例演示了我所描述的内容,但没有取消订阅功能、重复订阅检查或监督:

import scala.collection.mutable
import akka.actor.{Actor, ActorRef}

class HashBus() extends Actor {
  val topicActors = mutable.Map.empty[String, ActorRef]

  def createDistributionActor = {
    context.actorOf(Props[DistributionActor])
  }

  override def receive = {
    case subscribe : Subscribe =>
      topicActors.getOrElseUpdate(subscribe.topic, createDistributionActor) ! subscribe

    case publish : Publish =>
      topicActors.get(topic).foreach(_ ! publish)
  }
}

class DistributionActor extends Actor {

  val recipients = mutable.List.empty[ActorRef]

  override def receive = {
    case Subscribe(topic: String, actorRef: ActorRef) =>
      recipients +: actorRef

    case publish : Publish =>
      recipients.map(_ ! publish)
  }
}

这将确保您的总线 Actor 的邮箱不会饱和,因为总线的工作只是进行哈希查找。 DistributionActor 将负责映射接收者并分发有效负载。类似地,DistributionActor 可以保留主题的任何状态。

【讨论】:

  • 感谢您的建议!我正在考虑这个解决方案。但是存在一件事:主题解析规则更复杂并且具有掩码,这使得哈希解决方案更加复杂(查看它的规范github.com/butaji/JetMQ/blob/…)我很难解决的一件事 - 已解决订阅者的粘性:通常解决订阅者该主题比(取消)订阅本身更频繁的操作。这意味着我可以存储已处理的订阅和主题之间的链接
  • @VitalyBaum 您始终可以编写自己的哈希函数,将主题映射到唯一键(基于复杂的规则和掩码)。然后可以在 HashMap 中使用此哈希函数...
  • 我会试试这个解决方案。真的很喜欢分桶订阅的想法,这使得解决更快
  • 我在zeromq.org/whitepapers:message-matching987654322@这个话题上发现了相当有趣的文章
猜你喜欢
  • 1970-01-01
  • 2023-04-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-01-07
  • 2010-12-20
  • 1970-01-01
  • 2012-07-20
相关资源
最近更新 更多