Akka Cluster Singleton 完全指南:从 Classic 到 Typed API 的单例管理、故障切换与租约保障
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.mdClassic API 文档及其姊妹篇 akka-docs/src/main/paradox/typed/cluster-singleton.mdTyped 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 groupIdcom.typesafe.akka/groupId artifactIdakka-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/singletonsingleton-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-namesingleton子单例 Actor 的名称。role单例所在节点角色未指定则为全集群单例。hand-over-retry-interval1s新最老节点向可能正在离开的旧最老节点发送交接请求的重试间隔直到旧节点确认交接开始或旧节点被移除含akka.cluster.down-removal-margin。min-number-of-hand-over-retries15最小交接重试次数。实际重试次数由hand-over-retry-interval与akka.cluster.down-removal-margin推导但不少于该值。重试耗尽仍无法交换交接消息时管理器会抛出ClusterSingletonManagerIsStuck重启以恢复干净状态且仍不会启动单例直到旧最老节点被移出集群旧节点一侧则使用「重试次数 - 3」作为阈值之后停止单例实例。大集群可能需调大此值以避免 Leaving 到 Exiting 阶段的 gossip 传播导致过早超时正常退出场景下调小它不会让交接更快但极端故障下恢复可能更快。use-lease创建单例前要获取的租约配置路径租约丢失时 Actor 会重启并重新获取默认为无租约。lease-retry-interval5s获取租约的重试间隔。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-interval1s代理尝试解析单例实例的间隔。buffer-size1000单例位置未知时缓冲的消息数缓冲区满时新消息到达会丢弃最旧消息设为 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.scalacontext.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 pathowner 设为Cluster(system).selfAddress.hostPort。注意akka.cluster.singleton.lease-name配置键在此场景下被忽略。为单个单例配置租约可在配置中定义专属块LeaseDocSpec.scalamy.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),仅供参考