【发布时间】: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 负责存储订阅和主题的所有状态(例如保留消息)。
设备参与者可以决定他们想要发布、订阅或取消订阅主题的任何内容。
经过一些性能基准测试后,我意识到我当前的设计会影响发布和订阅之间的处理时间,原因如下:
- 我的 EventBus 实际上是一个单例
- 它造成了巨大的处理负载队列
如何为我的事件总线实施分配(并行化)工作负载? 当前的解决方案是否适合 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