【问题标题】:Broadcast Message to Routees in a ClusterRouter in Akka在 Akka 的 ClusterRouter 中向路由广播消息
【发布时间】:2013-05-20 21:47:24
【问题描述】:

我正在尝试向ClusterRouter 配置中的所有路由广播消息。我已经尝试了两种选择。这个:

 val workerRouter = context.actorOf(Props[ClusterRouter].withRouter(
    ClusterRouterConfig(AdaptiveLoadBalancingRouter(metrics), ClusterRouterSettings(
      totalInstances = 100, routeesPath = "/user/slave",
      allowLocalRoutees = true, useRole = None))), name = "slaveRouter")

  context.system.scheduler.schedule(2 seconds, 5 seconds, workerRouter, Broadcast(CapabilityRequest))

还有这个:

 val broadcastRouter = context.actorOf(Props[ClusterRouter].withRouter(
    ClusterRouterConfig(BroadcastRouter(Nil), ClusterRouterSettings(
      totalInstances = 100, routeesPath = "/user/slave",
      allowLocalRoutees = true, useRole = None))), name = "slaveRouter")

  context.system.scheduler.schedule(2 seconds, 5 seconds, broadcastRouter, CapabilityRequest)

但是对于他们两个来说,只有一个slaves 收到了消息。想法?


为了理解为什么我认为第一次尝试应该成功,必须查看AdaptiveLoadBalancingRouterLike trait 中的AdaptiveLoadBalancingRounter.scala,当Route 被创建时:

{
  case (sender, message) ⇒
    message match {
      case Broadcast(msg) ⇒ toAll(sender, routeeProvider.routees)
      case msg            ⇒ List(Destination(sender, getNext()))
    }
}

【问题讨论】:

  • 正如我在邮件列表中所问的那样:您能否在发送消息时证明您的集群实际上有多个成员?
  • 我可以给你整个代码,是的,但话又说回来:(i)RoundRobinRouter 向每个成员发送消息,(ii)“手动”广播,我在其中迭代所有成员网络似乎工作。

标签: scala akka broadcast akka-cluster


【解决方案1】:

在您的第一个示例中,您使用的路由器只会发送到一个路由。从我读过的文档中,该路由器将使用来自不同节点的可用指标来选择似乎受到最少胁迫的节点,并向该节点上的路由发送消息。我认为您在此设置中看到的行为是预期的。

对于您的第二个示例,我在文档中没有看到任何关于在集群环境中使用 BraodcastRouter 的内容,因此我不确定是否支持这种方法。话虽如此,我的猜测是使用空的路由列表(Nil)创建BraodcastRouter 是导致您看到的行为的原因。我认为如果您将其更改为BroadcastRouter(100),您可能会看到不同的行为。但同样,我不认为(基于文档中缺少示例)使用 BroadcastRouter 是受支持的(我可能是错的)。

您能否详细解释一下您的用例,以便我了解您的集群需要广播类型路由器的原因?

编辑

FWIW,我得到了使用以下代码的东西。一、配置:

akka {
  actor {
    provider = "akka.cluster.ClusterActorRefProvider"
  }
  remote {
    transport = "akka.remote.netty.NettyRemoteTransport"
    log-remote-lifecycle-events = off
    netty {
      hostname = "127.0.0.1"
      port = 0
    }
  }

  cluster {
    min-nr-of-members = 2
    seed-nodes = [
      "akka://ClusterSystem@127.0.0.1:2551", 
      "akka://ClusterSystem@127.0.0.1:2552"]

    auto-down = on
  }
}

然后,我使用以下代码启动了两个节点(一个在 2551,另一个在 2552):

object ClusterNode {

  def main(args: Array[String]): Unit = {

    // Override the configuration of the port 
    // when specified as program argument
    if (args.nonEmpty) System.setProperty("akka.remote.netty.port", args(0))


    // Create an Akka system
    val system = ActorSystem("ClusterSystem")
    val clusterListener = system.actorOf(Props(new Actor with ActorLogging {
      def receive = {
        case state: CurrentClusterState =>
          log.info("Current members: {}", state.members)
        case MemberJoined(member) =>
          log.info("Member joined: {}", member)
        case MemberUp(member) =>
          log.info("Member is Up: {}", member)
        case UnreachableMember(member) =>
          log.info("Member detected as unreachable: {}", member)
        case _: ClusterDomainEvent => // ignore

      }
    }), name = "clusterListener")

    Cluster(system).subscribe(clusterListener, classOf[ClusterDomainEvent])    
  }

}

class FooActor extends Actor{

  override def preStart = {
    println("Foo actor started on path: " + context.self.path)
  }

  def receive = {
    case msg => println(context.self.path + " received message: " + msg)
  }
}

然后我使用以下代码启动了第三个“节点”,即我的客户端节点:

object ClusterClient {
  def main(args: Array[String]) {
    val system = ActorSystem("ClusterSystem")

    Cluster(system) registerOnMemberUp{
      val router = system.actorOf(Props[FooActor].withRouter(
        ClusterRouterConfig(AdaptiveLoadBalancingRouter(HeapMetricsSelector),
        ClusterRouterSettings(
        totalInstances = 20, maxInstancesPerNode = 10,
        allowLocalRoutees = false))),
        name = "fooRouter")  

     router ! Broadcast("bar")
    }
  }
}

消息发送后,我看到它在两个服务器节点虚拟机中都收到了,每个虚拟机有 10 个参与者。

我的路由器和你的路由器之间的区别在于我没有指定本地路由,而是将routeesPath 换成了maxInstancesPerNode。我希望这会有所帮助。

【讨论】:

  • 有趣。我查看了所有源代码,但我也没有很好地解释为什么这不起作用。您是否看到在发送广播消息之前创建了 100 个从站?
  • 我将不得不通过你的例子,因为在我的情况下, routeesPath 需要理清哪些 routees 是合格的。我假设您在不同的虚拟机中运行不同的演员,对吧?
  • 是的。总而言之,我总共运行了 3 个 JVM(两个服务节点,一个客户端节点),但一切都在我的 mac 本地。在我的示例中,在我启动客户端代码并启动该集群感知路由器之前,不会部署参与者实例。我认为当服务节点启动时您自己在每个服务节点上启动参与者实例时情况会有所不同。
  • 我也在使用 Mac,所以...嗯...您使用的是哪个版本的 Akka? 2.1.x 还是 2.2-Mx ?
  • 好吧,也许我正在使用 2.2 的里程碑版本这一事实解释了这一点:P
猜你喜欢
  • 1970-01-01
  • 2020-12-09
  • 1970-01-01
  • 1970-01-01
  • 2016-12-02
  • 2017-05-19
  • 2013-09-09
  • 2020-11-23
  • 1970-01-01
相关资源
最近更新 更多