Akka Cluster Singleton 完全指南:从 Classic 到 Typed API 的单例管理、故障切换与租约保障
【免费下载链接】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 的 Cluster Singleton 模式正是为此设计:它通过ClusterSingletonManager与ClusterSingletonProxy两大组件,保证某个 Actor 在集群(或指定角色的节点组)中最多同时运行一份,并自动处理最老节点退出、崩溃时的优雅交接与接管。本文以 akka-docs/src/main/paradox/cluster-singleton.md(Classic API 文档)及其姊妹篇 akka-docs/src/main/paradox/typed/cluster-singleton.md(Typed API 文档)为骨架,结合akka-cluster-tools模块的源码、配置与 multi-jvm 测试,完整讲解核心概念、消息缓冲机制、全部配置参数、监督策略、终止消息与 Lease 租约,帮助你正确评估并在项目中落地这一模式。
模块信息与依赖引入
Cluster Singleton 功能位于akka-cluster-tools模块(Classic API)与akka-cluster-typed模块(Typed API)。官方文档要求通过 Akka 的安全仓库(需要 token 化 URL)获取依赖,使用 BOM 统一版本管理:
// sbt libraryDependencies += "com.typesafe.akka" %% "akka-cluster-tools" % AkkaVersion<!-- Maven --> <dependency> <groupId>com.typesafe.akka</groupId> <artifactId>akka-cluster-tools_2.13</artifactId> <version>${akka.version}</version> </dependency>// Gradle implementation 'com.typesafe.akka:akka-cluster-tools_2.13:${akka.version}'Typed API 则将 artifact 换成akka-cluster-typed_$scala.binary.version$。无论哪种 API,都需要集群已正确配置(akka.actor.provider = "cluster")。
核心概念:两个 Actor 与一条职责链
Cluster Singleton 模式由两个协作的 Actor 组成:
akka.cluster.singleton.ClusterSingletonManager:单例管理者。它必须在所有节点(或所有带指定角色的节点)上尽早启动。真正的单例 Actor 由它在最老节点上通过用户提供的Props创建为子 Actor。它保证任意时刻至多有一个单例实例在运行。akka.cluster.singleton.ClusterSingletonProxy:单例代理。它负责把消息路由到当前的单例实例。代理会持续跟踪集群中的最老节点,通过向单例的actorSelection显式发送akka.actor.Identify消息、等待回执来解析单例的ActorRef;若单例在配置时间内未回复,则周期性地重新解析。
最老节点由akka.cluster.Member#isOlderThan判定。当最老节点退出(Leaving)集群时,会先与新的最老节点执行一次交接(hand-over),之后才在新节点启动新的单例;因此在交接过程中会存在一段没有活跃单例的短暂窗口。若最老节点因 JVM 崩溃、强制关机或网络故障而不可达,集群的 failure detector 会发现异常,在节点被 Downing 并移除后,新的最老节点接管并创建新的单例——这种故障场景没有优雅交接,Akka 会尽力阻止出现多个活跃单例,极少数边界情况由可配置的超时最终解决,如需额外保险可叠加 Lease。
从 ClusterSingletonManager.scala 的源码结构看,管理器本身是一个基于 FSM(有限状态机)实现的 Actor,其状态包括Start、BecomingOldest、Oldest、WasOldest、HandingOver等,内部通过HandOverToMe、HandOverInProgress、HandOverDone等内部消息完成交接握手,这正对应文档描述的「优雅交接」流程。
Typed API 的入口:ClusterSingleton.init
在 Typed API 中,ClusterSingleton.init承担了「启动管理器 + 返回代理」双重职责。对给定singletonName调用init会返回一个ActorRef,向它发送消息即可到达单例实例,无需关心单例当前运行在哪个节点。init可被重复调用:若本节点已有同名的单例管理器运行,则不会额外启动管理器,只返回指向代理的ActorRef。
import akka.cluster.typed.{ ClusterSingleton, SingletonActor } val singletonManager = ClusterSingleton(system) // 需要时启动,并返回指向命名单例的代理 val proxy: ActorRef[Counter.Command] = singletonManager.init( SingletonActor(Behaviors.supervise(Counter()).onFailureException, "GlobalCounter")) proxy ! Counter.Increment经典 API 实战:JMS 队列消费单例
文档给出了一个典型的真实场景:外部系统的单一入口。假设有一个 JMS 队列,严格要求只有一个消费者存在以保证消息按顺序处理。先在所有节点启动ClusterSingletonManager并传入单例 Actor 的Props:
// Scala system.actorOf( ClusterSingletonManager.props( singletonProps = Props(classOf[Consumer], queue, testActor), terminationMessage = End, settings = ClusterSingletonManagerSettings(system).withRole("worker")), name = "consumer")// Java final ClusterSingletonManagerSettings settings = ClusterSingletonManagerSettings.create(system).withRole("worker"); system.actorOf( ClusterSingletonManager.props( Props.create(Consumer.class, () -> new Consumer(queue, testActor)), TestSingletonMessages.end(), settings), "consumer");代码来自 ClusterSingletonManagerSpec.scala 与 ClusterSingletonManagerTest.java。这里通过withRole("worker")把单例限制在带worker角色的节点上;若不指定withRole,则所有节点(不区分角色)均可承载单例。
随后从任意集群节点通过代理访问单例:
// Scala val proxy = system.actorOf( ClusterSingletonProxy.props( singletonManagerPath = "/user/consumer", settings = ClusterSingletonProxySettings(system).withRole("worker")), name = "consumerProxy")// Java ActorRef proxy = system.actorOf( ClusterSingletonProxy.props("/user/consumer", proxySettings), "consumerProxy");注意代理传入的是管理器的路径(/user/consumer),而不是单例本身的路径——单例是管理器的子 Actor,路径为/user/consumer/singleton(singleton-name可配置)。
应用特定终止消息:优雅关闭资源
管理器停止单例前会发送terminationMessage,用于让单例关闭外部资源。文档强调:PoisonPill是完全可用的终止消息,但若需要先释放 JMS 连接等资源,则应使用应用自定义的消息:
// Scala —— 单例收到 End 后先注销消费者,收到 UnregistrationOk 再停止自身 case End => queue ! UnregisterConsumer case UnregistrationOk => stoppedBeforeUnregistration = false context.stop(self)该片段取自 ClusterSingletonManagerSpec.scala 中Consumer的实现:preStart时向队列注册,End触发注销流程,确认注销成功后才真正停止——通过「先注销、后停止」保证切换期间绝不会出现两个消费者同时连接队列。
Typed API 中对应的是withStopMessage,在 SingletonCompileOnlySpec.scala 中可见:SingletonActor(Counter(), "GlobalCounter").withStopMessage(Counter.GoodByeCounter)。Typed 文档补充了一个要点:交接到新最老节点的流程在单例 Actor 终止后才算完成;如果关闭逻辑不含异步操作,可以直接写在PostStop信号处理器中。
消息缓冲:代理的容错窗口
由于代理需要周期性地解析单例位置,在节点离开集群等场景下会出现ActorRef暂时不可用的窗口。此时代理会将发往单例的消息缓冲起来,待单例可用后投递;若缓冲区已满,新消息到达时会丢弃最旧的消息。缓冲大小可配置,设为 0 表示完全禁用缓冲(位置未知时立即丢弃消息)。
文档同时给出重要提醒:由于这些 Actor 的分布式本质,消息总是可能丢失,应在单例侧实现确认(acknowledgement)、在客户端实现重试(retry),以达成至少一次(at-least-once)投递语义。此外,单例不会运行在 WeaklyUp 状态的成员上。
配置详解:全部参数与默认值
以下配置块完整取自 reference.conf,即 typed/cluster-singleton.md 中引用的#singleton-config与#singleton-proxy-config两段。
ClusterSingletonManager 配置(akka.cluster.singleton)
| 配置键 | 默认值 | 说明 |
|---|---|---|
singleton-name | "singleton" | 子单例 Actor 的名称。 |
role | "" | 单例所在节点角色;未指定则为全集群单例。 |
hand-over-retry-interval | 1s | 新最老节点向可能正在离开的旧最老节点发送交接请求的重试间隔,直到旧节点确认交接开始,或旧节点被移除(含akka.cluster.down-removal-margin)。 |
min-number-of-hand-over-retries | 15 | 最小交接重试次数。实际重试次数由hand-over-retry-interval与akka.cluster.down-removal-margin推导,但不少于该值。重试耗尽仍无法交换交接消息时,管理器会抛出ClusterSingletonManagerIsStuck重启以恢复干净状态,且仍不会启动单例,直到旧最老节点被移出集群;旧节点一侧则使用「重试次数 - 3」作为阈值,之后停止单例实例。大集群可能需调大此值以避免 Leaving 到 Exiting 阶段的 gossip 传播导致过早超时;正常退出场景下调小它不会让交接更快,但极端故障下恢复可能更快。 |
use-lease | "" | 创建单例前要获取的租约配置路径;租约丢失时 Actor 会重启并重新获取;默认为无租约。 |
lease-retry-interval | 5s | 获取租约的重试间隔。 |
lease-name | "" | 自定义租约名。注意多个单例必须使用唯一租约名,可通过ClusterSingletonSettings的 leaseSettings 定义;未定义时由 ActorSystem 名与单例 Actor 路径推导,但可能过长。Typed 集群无法通过此配置修改,任何值都会被忽略,必须通过编程 API 设置。 |
从 ClusterSingletonManagerSettings.apply 的源码可以看到:role为空字符串会被转换为None(全集群),use-lease为空则返回None(无租约);removalMargin在默认构造中显式设为Duration.Zero,并注释说明实际会回退到DowningProvider.downRemovalMargin——这也是文档中「重试次数与 down-removal-margin 联动」的实现依据。manager 与 proxy 的设置均可通过withXxx方法按单例粒度定制。
ClusterSingletonProxy 配置(akka.cluster.singleton-proxy)
| 配置键 | 默认值 | 说明 |
|---|---|---|
singleton-name | ${akka.cluster.singleton.singleton-name} | 管理器启动的单例 Actor 名称,与 manager 侧保持一致。 |
role | "" | 单例可部署的节点角色,须与ClusterSingletonManager的角色一致;未指定则为全集群。 |
singleton-identification-interval | 1s | 代理尝试解析单例实例的间隔。 |
buffer-size | 1000 | 单例位置未知时缓冲的消息数;缓冲区满时新消息到达会丢弃最旧消息;设为 0 禁用缓冲(立即丢弃);最大允许 10000。 |
监督(Supervision):两个层级,两种策略
单例涉及两个可被监督的 Actor(以文中的consumer为例):
- 集群单例管理器,如
/user/consumer,运行在集群每个节点上; - 用户单例 Actor,如
/user/consumer/singleton,由管理器在最老节点上启动。
管理器不应修改监督策略——它必须始终运行。需要监督的是用户单例。文档给出的做法是引入一个父级监督 Actor,由它来创建「真正的」单例实例:
// Scala class SupervisorActor(childProps: Props, override val supervisorStrategy: SupervisorStrategy) extends Actor { val child = context.actorOf(childProps, "supervised-child") def receive = { case msg => child.forward(msg) } }使用时将SupervisorActor作为singletonProps,把真实单例的Props与自定义SupervisorStrategy传入(ClusterSingletonSupervision.scala):
context.system.actorOf( ClusterSingletonManager.props( singletonProps = Props(classOf[SupervisorActor], props, supervisorStrategy), terminationMessage = PoisonPill, settings = ClusterSingletonManagerSettings(context.system)), name = name)Typed API 的监督更直接:默认策略是异常抛出时停止 Actor,通过Behaviors.supervise(...).onFailureException覆盖为重启以保证常驻;也可以使用带退避的重启:
val proxyBackOff: ActorRef[Counter.Command] = singletonManager.init( SingletonActor( Behaviors .supervise(Counter()) .onFailureException), "GlobalCounter"))来自 SingletonCompileOnlySpec.scala。注意退避重启意味着存在单例暂不运行的窗口,完整的监督选项参见 fault-tolerance。
Lease 租约:防止双单例的最后保险
即使在合理配置下,仍存在同时出现两个单例的极少数可能:无合适 downing provider 的网络分区、部署失误导致两个独立 Akka 集群、分区两侧「移除成员」与「关闭节点」的时序差异。Lease(见 coordination)可作为最终备份——获取不到租约就不创建单例 Actor。
全局启用:在application.conf中设置akka.cluster.singleton.use-lease为所用租约的配置位置。租约名形如<actor system name>-singleton-<singleton actor path>,owner 设为Cluster(system).selfAddress.hostPort。注意akka.cluster.singleton.lease-name配置键在此场景下被忽略。
为单个单例配置租约,可在配置中定义专属块(LeaseDocSpec.scala):
my.app.my-singleton-lease { use-lease = "akka.coordination.lease.kubernetes" lease-retry-interval = 5s lease-name = "my-pingpong-singleton-lease" }然后从配置加载,或编程指定(LeaseDocSpec.scala):
// 从配置加载 val settings = ClusterSingletonSettings(system).withLeaseSettings( LeaseUsageSettings(system.settings.config.getConfig("my.app.my-singleton-lease"))) // 编程指定 val settings2 = ClusterSingletonSettings(system).withLeaseSettings( LeaseUsageSettings("akka.coordination.lease.kubernetes", 5.seconds, "my-pingpong-singleton-lease")) val singletonActor = SingletonActor(pingPong, "ping-pong").withStopMessage(Perish).withSettings(settings) ClusterSingleton(system).init(singletonActor)租约行为:管理器作为最老节点却获取不到租约时会持续重试;租约丢失则终止单例 Actor 后重新尝试获取(源码中对应LeaseLost事件与DelayedLeaseRetry重试机制,见 ClusterSingletonManager.scala)。
使用前必读:潜在问题与注意事项
文档明确警示该模式不应成为首选设计,它有明显代价:
- 性能瓶颈:单例可能迅速成为瓶颈点;
- 非零停机:不能依赖单例持续可用——承载单例的节点死亡后,需要数秒才能被检测到并迁移到其他节点;
- 最老节点集中:多个单例全部运行在最老节点(或指定角色的最老节点)上;若单例数量多,可考虑 Cluster Sharding 结合常驻实体作为更优替代。
最重要的警告:切勿使用可能把集群分裂成多个独立集群的 downing 策略(网络问题或长 GC 暂停时),否则每个分裂出的集群都会启动一个单例,出现多个单例并存!务必阅读 Downing 相关章节。
行为验证:multi-jvm 测试揭示的真实语义
akka-cluster-tools提供了完整的 multi-jvm 测试 ClusterSingletonManagerSpec.scala,用 6 个节点逐一验证了文档描述的语义:
- 6 节点(其中 5 个带
worker角色)启动后,单例注册并运行在第一个加入的最老节点上; - 代理可从任意节点把消息路由到最老节点上的单例(
verifyProxyMsg会断言回包确实来自最老节点的地址); - 最老节点调用
Cluster(system).leave(...)优雅退出后,单例完成交接并在新最老节点上重新注册(verifyRegistration(second)),旧节点的单例被终止; - 最老节点崩溃(
testConductor.exit模拟)后,新的最老节点接管并重新创建单例,连续验证「5 节点」、「3 节点」、「2 节点」三种故障规模下的接管。
测试中的PointToPointChannel是「极其严格」的点对点通道:任何「重复注册」或「非预期注销」都会导致通道自我终止,从而把「两个单例并存」的行为直接暴露为测试失败——这正是文档「至多一个单例实例」承诺的可执行验证。消息类与序列化(CborSerializable)的完整定义见 TestSingletonMessages.java。
小结
Cluster Singleton 是 Akka 中「恰好一个实例」问题的标准答案:ClusterSingletonManager用 FSM 保证单例唯一性,ClusterSingletonProxy用 Identify 探测与消息缓冲消化位置切换窗口,Lease 租约兜底极端故障场景。它适合单一协调点、外部系统单一入口、单主多从等场景,但也必须清醒认识其瓶颈、非零停机与最老节点集中等代价,并在 downing 策略上格外谨慎。文中所有配置均以 reference.conf 的实际默认值为准,源码路径与测试用例可在本仓库中进一步深入研读。
【免费下载链接】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
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考