ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Akka Classic 集群感知路由器(Cluster Aware Routers)完整指南:Group 与 Pool 的配置、原理与实战

Akka Classic 集群感知路由器(Cluster Aware Routers)完整指南:Group 与 Pool 的配置、原理与实战 后端并发编程异步编程【免费下载链接】akka-coreA platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.项目地址https://gitcode.com/gh_mirrors/ak/akka-core点击查看免费下载集群感知路由器Cluster Aware Router是 Akka Classic Actor API 中一类特殊的路由器它不再把 routee 固定部署在本地进程内而是把 routee 的部署与查找与集群成员member状态动态绑定——节点加入、离开、失联或恢复时routee 的集合会随之自动伸缩。本文将围绕 akka-docs/src/main/paradox/cluster-routing.md 展开系统讲解 Group 与 Pool 两种集群路由器的工作方式、完整的 HOCON 配置与代码定义方式、底层的事件驱动原理对应 ClusterRouterConfig.scala 的实现并结合 akka-cluster-metrics 模块中的 StatsService 分布式示例给出可复制、可运行的完整实战代码。读完本文你将能够独立配置并理解基于集群成员动态伸缩的 Actor 路由基础设施。说明本文面向Classic经典Actor API。Akka 官方建议新项目使用 Typed API其对应的路由器文档见 typed/routers.md但 Classic 集群路由器在存量项目中仍被广泛使用且其核心机制成员感知、角色过滤、routee 自动增删是理解 Akka 集群路由的基础。一、什么是集群感知路由器核心工作机制所有 Classic 路由器见 routing.md都可以“感知”集群中的成员节点即在集群节点上部署新的 routee或查找已有的 routee。其核心行为可以总结为三条自动规则节点失联unreachable或离开leave集群该节点上的 routee 会被自动从路由器注销unregister新节点加入集群路由器会根据配置自动为其添加 routee失联节点恢复可达reachable again该节点上的 routee 会被重新注册。此外如果集群启用了 WeaklyUp 特性集群感知路由器也会把处于WeaklyUp状态的成员纳入路由范围。两种集群路由器类型类型工作方式routee 归属典型场景Group通过 Actor 路径选择ActorSelection把消息发送到指定路径的 actorroutee 可被集群中不同节点上的多个路由器共享后端节点运行服务前端节点上的路由器调用它形成“服务发现式”路由Pool路由器把 routee 作为子 actor创建并部署远程部署到集群节点上每个路由器拥有自己的routee 实例互不共享单一 master 协调任务把实际工作委托给集群其他节点上的 routeePool 的伸缩示例在一个 10 节点的集群中如果在 3 个节点上启动了路由器并且配置为每节点一个实例那么总共会有 30 个 routee——每个路由器各创建 10 个它们之间互不共享。关键差异Group 的 routee 由应用自行启动路由器只负责“找到它们”Pool 的 routee 由路由器负责创建和部署应用不直接管理。二、依赖准备集群感知路由器需要akka-cluster模块。Akka 的依赖从 Akka 官方安全库仓库获取需使用带 token 的安全 URL参见 Akka 账户页面说明。以 sbt 为例// 使用 akka-bom 统一管理版本AkkaVersion 由 BOM 符号替换 libraryDependencies com.typesafe.akka %% akka-cluster % AkkaVersionMaven 与 Gradle 的对应坐标均为com.typesafe.akka:akka-cluster_scalaBinaryVersion:version实际依赖声明可从项目 artifact-bom 中获得版本对齐建议。三、Group路由到一组集群成员上的既有 actor3.1 使用前提与配置使用Group时routee actor 必须由应用自己或应用的其他机制在集群成员节点上启动路由器不会帮你创建它们。典型的 HOCON 配置如下akka.actor.deployment { /statsService/workerRouter { router consistent-hashing-group routees.paths [/user/statsWorker] cluster { enabled on allow-local-routees on use-roles [compute] } } }3.2 关键配置项详解配置项含义默认值来自 reference.confrouter底层路由逻辑如consistent-hashing-group、round-robin-group等无routees.pathsroutee 的 actor 路径列表不能包含协议与地址信息如akka://...host:port地址信息由路由器从集群成员关系中动态获取无cluster.enabled是否启用集群感知offcluster.max-total-nr-of-instances集群中 routee 的总数上限10000cluster.max-nr-of-instances-per-node每个节点上的 routee 数上限Group 场景通常为 11cluster.allow-local-routees是否允许 routee 与本路由器位于同一节点oncluster.use-roles仅使用带有指定角色集合的成员节点未定义或为空时使用所有成员[]几点需要特别注意的行为路径不含地址routees.paths只写 actor 的相对路径如/user/statsWorker因为节点地址是运行期从集群成员动态获取的转发语义消息通过 ActorSelection 转发给 routee因此投递语义与直接使用 ActorSelection 一致max-total-nr-of-instances默认值极大10000这意味着只要新节点加入路由器就会持续为其添加 routee如果你希望限制 routee 总量必须显式调低该值角色过滤通过use-roles可以把 routee 查找限制在带有特定角色的成员节点上如只路由到compute角色的节点启动时机routee actor 应尽可能早地在 actor system 启动时创建因为一旦成员状态变为Up路由器会立刻尝试使用这些 routee。这一点在文档中特别以 note 形式强调。兼容性细节从源码 ClusterRouterSettingsBase.getMaxTotalNrOfInstances 可以看到出于向后兼容nr-of-instances与max-total-nr-of-instances在集群感知路由器中语义相同如果用户显式定义了nr-of-instances且非 0/1它优先于max-total-nr-of-instances。另外 ClusterRouterGroupSettings.fromConfig 会把已废弃的cluster.use-roleAkka 2.5.4 起被use-roles取代与cluster.use-roles合并处理。3.3 在代码中定义 Group 路由器除了 HOCON 配置也可以在代码中直接构造ClusterRouterGroup完整示例见 StatsService.scala 与 StatsService.javaScalaimport akka.cluster.routing.{ ClusterRouterGroup, ClusterRouterGroupSettings } import akka.routing.ConsistentHashingGroup val workerRouter context.actorOf( ClusterRouterGroup( ConsistentHashingGroup(Nil), ClusterRouterGroupSettings( totalInstances 100, routeesPaths List(/user/statsWorker), allowLocalRoutees true, useRoles Set(compute))).props(), name workerRouter2)Javaint totalInstances 100; IterableString routeesPaths Collections.singletonList(/user/statsWorker); boolean allowLocalRoutees true; SetString useRoles new HashSet(Arrays.asList(compute)); ActorRef workerRouter getContext() .actorOf( new ClusterRouterGroup( new ConsistentHashingGroup(routeesPaths), new ClusterRouterGroupSettings( totalInstances, routeesPaths, allowLocalRoutees, useRoles)) .props(), workerRouter2);从 ClusterRouterGroupSettings 的构造约束可以看到totalInstances必须大于 0routeesPaths不能为空且每个路径必须是不含地址信息的相对 actor 路径否则会抛出IllegalArgumentException。3.4 Group 内部原理如何“找到”routee从 ClusterRouterGroupActor 的实现可以看到 Group 路由器的运行机制路由器 actor 在preStart时向Cluster订阅MemberEvent与ReachabilityEvent见 ClusterRouterActor.preStartselectDeploymentTarget每次选出一个「节点 路径」组合优先选择尚未使用过的节点其次选择已用路径最少的节点对每个选中目标通过group.routeeFor(address.toString path, context)构造一个以节点地址为 anchor 的 ActorSelection routee当成员状态变化Up/WeaklyUp、失联、离开时clusterReceive会相应调用addMember/removeMember从而动态增删 routee。关键的一点isAvailable要求成员状态为Up或WeaklyUp、满足角色过滤且当allowLocalRoutees off时排除本节点。四、Pool把 routee 远程部署到集群节点4.1 配置示例使用Pool时routee 由路由器作为子 actor 创建并远程部署到集群成员节点上。配置如下akka.actor.deployment { /statsService/singleton/workerRouter { router consistent-hashing-pool cluster { enabled on max-nr-of-instances-per-node 3 allow-local-routees on use-roles [compute] } } }4.2 与 Group 的配置差异Pool 与 Group 共享cluster.enabled、cluster.max-total-nr-of-instances、cluster.allow-local-routees、cluster.use-roles等配置区别在于Pool 不配置routees.pathsroutee 不是查找出来的而是创建出来的Pool 通常显式设置max-nr-of-instances-per-node每节点上限上面示例为 3max-total-nr-of-instances仍然定义集群中 routee 的总上限但每节点的 routee 数不会超过max-nr-of-instances-per-node。例如设置max-total-nr-of-instances 50、max-nr-of-instances-per-node 2时路由器会为每个新加入的成员部署 2 个 routee直到达到 25 个成员或总量上限reference.conf 中的注释正是这一示例。4.3 在代码中定义 Pool 路由器Scalaimport akka.cluster.routing.{ ClusterRouterPool, ClusterRouterPoolSettings } import akka.routing.ConsistentHashingPool val workerRouter context.actorOf( ClusterRouterPool( ConsistentHashingPool(0), ClusterRouterPoolSettings(totalInstances 100, maxInstancesPerNode 3, allowLocalRoutees false)) .props(Props[StatsWorker]()), name workerRouter3)Javaint totalInstances 100; int maxInstancesPerNode 3; boolean allowLocalRoutees false; SetString useRoles new HashSet(Arrays.asList(compute)); ActorRef workerRouter getContext() .actorOf( new ClusterRouterPool( new ConsistentHashingPool(0), new ClusterRouterPoolSettings( totalInstances, maxInstancesPerNode, allowLocalRoutees, useRoles)) .props(Props.create(StatsWorker.class)), workerRouter3);4.4 Pool 的部署策略最少负载节点优先ClusterRouterPoolActor.selectDeploymentTarget 展示了 Pool 的部署算法每次从当前可用节点中选出routee 数量最少的节点只要该节点上的 routee 数未超过maxInstancesPerNode且总数未达totalInstances就通过RemoteScope(target)部署一个新的 routee见 ClusterRouterPoolActor.addRoutees。4.5 使用 Pool 的序列化前提使用远程部署的 Pool 时Props的所有参数都必须可序列化参见 serialization.md否则远程部署会失败。这是文档明确强调的前提条件。4.6 Pool 的限制从源码看ClusterRouterPool有一个明确约束不允许同时使用 Resizer动态调整器构造时若local.resizer非空会直接抛出异常ClusterRouterPool。因为集群感知路由器本身已经具备基于成员变化的动态伸缩能力再叠加 Resizer 没有意义。五、实战示例一Group 路由——分布式文本统计服务StatsService文档以「统计文本中单词平均长度」的服务为例完整演示了 Group 路由器的用法。全部源码位于Scalaakka-cluster-metrics/src/multi-jvm/scala/akka/cluster/metrics/sample/下的 StatsService.scala、StatsWorker.scala、StatsMessages.scalaJavaakka-docs/src/test/java/jdocs/cluster/下的 StatsService.java、StatsAggregator.java、StatsMessages.java5.1 业务逻辑概述整个服务的流程是收到一段文本 → 拆分为单词 → 把「统计每个单词字符数」的任务委托给 worker即 router 的 routee→ worker 返回字符数 → 聚合器aggregator收集全部结果并计算「平均每词字符数」。5.2 消息定义Scala使用CborSerializable以便跨节点序列化final case class StatsJob(text: String) extends CborSerializable final case class StatsResult(meanWordLength: Double) extends CborSerializable final case class JobFailed(reason: String) extends CborSerializableJava实现Serializablepublic interface StatsMessages { public static class StatsJob implements Serializable { private final String text; public StatsJob(String text) { this.text text; } public String getText() { return text; } } public static class StatsResult implements Serializable { /* meanWordLength: Double */ } public static class JobFailed implements Serializable { /* reason: String */ } }5.3 Worker单词字符数统计带缓存ScalaStatsWorker.scalaclass StatsWorker extends Actor { var cache Map.empty[String, Int] def receive { case word: String val length cache.get(word) match { case Some(x) x case None val x word.length cache (word - x) x } sender() ! length } }Java版本逻辑相同对收到的String单词先在本地缓存中查长度未命中则计算word.length()并写入缓存最后把长度回发给sender()。5.4 Service拆分单词并委托给路由器ScalaStatsService.scalaclass StatsService extends Actor { // 该路由器同时用于 lookup 与 deploy 两种场景 // 若只用 lookupGroup可把 Props[StatsWorker]() 换成 Props.empty val workerRouter context.actorOf(FromConfig.props(Props[StatsWorker]()), name workerRouter) def receive { case StatsJob(text) if text ! val words text.split( ) val replyTo sender() // 重要不能在闭包中捕获 sender() // 创建聚合 actor 收集 worker 的回复 val aggregator context.actorOf(Props(classOf[StatsAggregator], words.size, replyTo)) words.foreach { word workerRouter.tell(ConsistentHashableEnvelope(word, word), aggregator) } } }JavaStatsService.java逻辑一致FromConfig.getInstance().props(Props.create(StatsWorker.class))创建workerRouter收到非空StatsJob后拆词、创建StatsAggregator并用ConsistentHashableEnvelope(word, word)把每个单词按内容哈希路由到 worker。注意两点工程细节其一replyTo sender()必须在闭包外用局部变量保存否则闭包内读取到的可能是错误的 sender其二这里使用ConsistentHashableEnvelope保证同一个单词总是被路由到同一个 worker从而充分利用 worker 内部的缓存这是选择consistent-hashing路由策略的原因。5.5 Aggregator聚合结果或超时失败ScalaStatsService.scalaclass StatsAggregator(expectedResults: Int, replyTo: ActorRef) extends Actor { var results IndexedSeq.empty[Int] context.setReceiveTimeout(3.seconds) def receive { case wordCount: Int results results : wordCount if (results.size expectedResults) { val meanWordLength results.sum.toDouble / results.size replyTo ! StatsResult(meanWordLength) context.stop(self) } case ReceiveTimeout replyTo ! JobFailed(Service unavailable, try again later) context.stop(self) } }JavaStatsAggregator.java在preStart中设置 3 秒的ReceiveTimeout收齐expectedResults个结果就计算平均值并回发StatsResult超过 3 秒未收齐则回发JobFailed(Service unavailable, try again later)两种情况都主动stop自己避免泄漏。5.6 部署与路由配置所有节点都启动StatsService与StatsWorker注意Group 场景下 worker 是 routee必须由各节点自行启动。路由器的 HOCON 配置为akka.actor.deployment { /statsService/workerRouter { router consistent-hashing-group routees.paths [/user/statsWorker] cluster { enabled on allow-local-routees on use-roles [compute] } } }这样用户请求可以发送到任意节点上的StatsService而它会通过集群路由器把单词统计任务分发到所有节点上的StatsWorker。整个示例中业务 actor消息、worker、service、aggregator本身不包含任何集群相关代码——集群感知能力完全由路由器配置注入这正是该设计的精髓。六、实战示例二Pool 路由——单 Master 远程部署 worker6.1 场景与思路第二个示例演示「单 master 创建并远程部署 worker」集群中只有一个StatsService单例作为任务协调者它通过 Pool 路由器把 worker 部署到集群各节点上。为了保证「只有一个 master」示例借助akka-cluster-tools模块的Cluster Singleton见 cluster-singleton.md。6.2 启动 ClusterSingletonManager每节点都启动Scalasystem.actorOf( ClusterSingletonManager.props( singletonProps Props[StatsService], terminationMessage PoisonPill, settings ClusterSingletonManagerSettings(system).withRole(compute)), name statsService)JavaStatsSampleOneMasterMain.javaClusterSingletonManagerSettings settings ClusterSingletonManagerSettings.create(system).withRole(compute); system.actorOf( ClusterSingletonManager.props( Props.create(StatsService.class), PoisonPill.getInstance(), settings), statsService);ClusterSingletonManager会在每个节点上启动但只有「最老」的节点满足角色要求真正运行StatsService单例。6.3 启动 ClusterSingletonProxy每节点都启动Scalasystem.actorOf( ClusterSingletonProxy.props( singletonManagerPath /user/statsService, settings ClusterSingletonProxySettings(system).withRole(compute)), name statsServiceProxy)JavaStatsSampleOneMasterMain.javaClusterSingletonProxySettings proxySettings ClusterSingletonProxySettings.create(system).withRole(compute); system.actorOf( ClusterSingletonProxy.props(/user/statsService, proxySettings), statsServiceProxy);ClusterSingletonProxy接收用户的文本请求通过监听集群事件找到当前StatsService运行在最老节点上的单例 master并把任务委托给它。6.4 Pool 路由器配置akka.actor.deployment { /statsService/singleton/workerRouter { router consistent-hashing-pool cluster { enabled on max-nr-of-instances-per-node 3 allow-local-routees on use-roles [compute] } } }这里路径statsService/singleton/workerRouter中多出的singleton段是因为ClusterSingletonManager会在其内部创建名为singleton的子 actor 来承载单例——因此单例内部的workerRouter的完整部署路径是/user/statsService/singleton/workerRouter。6.5 运行方式文档提示最容易的动手方式是直接运行官方集群示例包包含如何运行Router Example with Pool of Routees的说明对应仓库内的示例工程为 samples/akka-sample-cluster-scala 与 samples/akka-sample-cluster-java。Java 版本的启动入口 StatsSampleOneMasterMain.java 演示了通过命令行参数传入端口启动多节点如2551 2552 0其中0表示随机端口并用ConfigFactory.parseString(akka.cluster.roles [compute])为节点标记角色、以ConfigFactory.load(stats2)加载集群配置——这正是use-roles [compute]生效的前提。七、参数与行为速查表以下汇总集群感知路由器的全部核心配置及其含义默认值来自 reference.conf配置路径类型默认值说明akka.actor.deployment.path.router字符串无底层路由策略Group 用*-group如consistent-hashing-groupPool 用*-pool如consistent-hashing-poolroutees.paths字符串数组无仅 Grouproutee 相对路径列表不含协议与地址cluster.enabled布尔off是否启用集群感知cluster.max-total-nr-of-instances整数10000集群中 routee 总上限兼容别名nr-of-instances优先cluster.max-nr-of-instances-per-node整数1每节点 routee 上限cluster.allow-local-routees布尔on是否允许本地 routeemaster-worker 场景可设off强制全远程cluster.use-roles字符串数组[]仅使用带全部指定角色的成员空则使用所有成员cluster.use-role字符串已废弃Akka 2.5.4 起由use-roles取代源码中仍做兼容合并行为速查节点Up/WeaklyUp→ 添加 routee节点失联 /Exited/Removed→ 注销该节点 routee失联恢复可达 → 重新添加max-total-nr-of-instances默认 10000意味着节点加入就会持续扩容需主动调低以限制总量Group 的 routee 通过 ActorSelection 转发投递语义与非集群 ActorSelection 一致Pool 部署时选择「当前 routee 最少」的节点且每节点不超过max-nr-of-instances-per-nodePool 的Props参数必须可序列化Pool 不可与 Resizer 同时使用。八、从源码理解集群事件如何驱动 routee 增删最后我们把文档描述与底层实现对应起来实现位于 ClusterRouterConfig.scala订阅路由器 actor 在preStart中向Cluster订阅MemberEvent与ReachabilityEventL528-L531并在postStop退订初始状态以cluster.readView.members中所有可用成员Up/WeaklyUp且满足角色的地址初始化节点集合L533-L538事件分发clusterReceive处理CurrentClusterState、MemberEvent、UnreachableMember、ReachableMemberL594-L613——可用成员事件触发addMember节点加入集合并addRoutees不可用/失联事件触发removeMember节点移出集合、注销该节点的 routee、必要时补建新 routee伸缩约束addRoutees遵循totalInstances与maxInstancesPerNode双重上限selectDeploymentTarget中currentRoutees.size totalInstances或节点count maxInstancesPerNode时停止创建可用性判定isAvailable要求状态为Up/WeaklyUp、角色满足useRoles.subsetOf(memberRoles)、且allowLocalRoutees off时排除本节点L540-L545。这套事件驱动机制保证了路由器在集群拓扑变化时无需人工干预即可自愈是 Akka 弹性elasticity与韧性resilience在消息路由层的直接体现。九、进一步阅读路由器通用概念Classicrouting.md集群感知路由器Typed API新项目推荐typed/routers.md集群成员状态与 WeaklyUptyped/cluster-membership.md集群单例单 master 场景依赖cluster-singleton.md完整配置参考general/configuration-reference.md可运行的官方示例Scala 版 samples/akka-sample-cluster-scala、Java 版 samples/akka-sample-cluster-java赞分享后端并发编程异步编程【免费下载链接】akka-coreA platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.项目地址https://gitcode.com/gh_mirrors/ak/akka-core点击查看免费下载相关推荐Akka Classic Routing 全面指南路由器的 7 种路由策略、Pool/Group 形态与自定义实现Akka Classic Routing 全面指南路由器的 7 种路由策略、Pool/Group 形态与自定义实现 导读 Akka Classic Routi后端并发编程异步编程Akka Classic Cluster 实战指南集群扩展、成员事件订阅、种子节点加入与运维配置Akka Classic Cluster 实战指南集群扩展、成员事件订阅、种子节点加入与运维配置 Akka Classic Cluster 是基于经典 Act后端并发编程异步编程AIBrix 前缀缓存感知路由Prefix Cache Aware Routing全解析原理、配置与 KV 事件同步实战AIBrix 前缀缓存感知路由Prefix Cache Aware Routing全解析原理、配置与 KV 事件同步实战 Prefix Cache Awa云原生大模型模型推理服务API网关LLM 网关弹性伸缩可观测性后端创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表