ARTICLE DETAIL

资讯详情

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

Akka Classic Cluster Sharding 完全指南:实体分片、状态存储与优雅运维

Akka Classic Cluster Sharding 完全指南:实体分片、状态存储与优雅运维 Akka Classic Cluster Sharding 完全指南实体分片、状态存储与优雅运维【免费下载链接】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-coreCluster Sharding集群分片是 Akka 中用于在集群多个节点间自动分布有状态 Actor称为实体 Entity的核心扩展应用只需通过逻辑实体标识符与实体交互无需关心其物理位置。本文以经典ClassicActor API 下的 cluster-sharding.md 文档为主体结合仓库中的多节点测试与配置源码系统讲解从基本用法、分片算法设计、状态存储模式、Passivation 到优雅停机与状态巡检的完整实战方案。读完本文你将掌握如何在 Akka Cluster 中注册实体类型、编写消息提取器、选择 Coordinator 状态存储模式并对分片状态进行监控与治理。说明本文聚焦 Classic Actor API。Akka 官方建议新项目使用 Typed API其完整文档见 typed/cluster-sharding.md经典 API 的完整功能文档也以该 Typed 文档为准两者共享大量概念。模块依赖要使用 Cluster Sharding需要引入akka-cluster-sharding依赖基于 Akka BOM 管理版本// sbt libraryDependencies com.typesafe.akka %% akka-cluster-sharding % AkkaVersion!-- Maven -- dependency groupIdcom.typesafe.akka/groupId artifactIdakka-cluster-sharding_2.13/artifactId version${akka.version}/version /dependencyAkka 依赖托管在 Akka 的安全库仓库secure library repository中需要按照 https://account.akka.io/token 的说明使用带 token 的安全 URL 访问。该模块对应的源码位于 akka-cluster-sharding其内部消息的 Protobuf 定义见 ClusterShardingMessages.proto。引入与核心概念Cluster Sharding 解决的核心问题是当大量有状态 Actor 的总资源消耗如内存超过单机承载能力时将 Actor 按标识符自动分布到集群多个节点。每个实体 Actor 在同一时刻只在一个节点运行发送方无需知道其物理位置——消息统一经本地的ShardRegionActor 路由由扩展负责查找实体所在位置。本文对应的实现与验证代码主要有两处Scala 多节点测试ClusterShardingSpec.scalaJava 文档示例ClusterShardingTest.java两个文件均包含本文引用的#counter-actor、#counter-start、#counter-extractor、#counter-usage等标记片段是理解下文示例的第一手材料。基本示例一个计数器实体1. 定义实体 Actor以下是一个基于事件溯源Event Sourcing的持久化计数器实体// 来源ClusterShardingSpec.scala #counter-actor case object Increment case object Decrement final case class Get(counterId: Long) final case class EntityEnvelope(id: Long, payload: Any) case object Stop final case class CounterChanged(delta: Int) class Counter extends PersistentActor { import ShardRegion.Passivate context.setReceiveTimeout(120.seconds) // self.path.name 即实体标识符utf-8 URL 编码 override def persistenceId: String Counter- self.path.name var count 0 def updateState(event: CounterChanged): Unit count event.delta override def receiveRecover: Receive { case evt: CounterChanged updateState(evt) } override def receiveCommand: Receive { case Increment persist(CounterChanged(1))(updateState) case Decrement persist(CounterChanged(-1))(updateState) case Get(_) sender() ! count case ReceiveTimeout context.parent ! Passivate(stopMessage Stop) case Stop context.stop(self) } }对应 Java 版本AbstractPersistentActor见 ClusterShardingTest.java 的#counter-actor片段。两个关键点实体不必须是持久化 Actor。文档明确指出如果实体在节点故障或迁移后其状态仍有价值就必须能够恢复状态但 Cluster Sharding 本身并不强制持久化。persistenceId的构造方式self.path.name就是实体标识符utf-8 URL 编码因此Counter- self.path.name能保证每个实体拥有唯一持久化 ID。你可以用其他方式定义但必须保证唯一。2. 注册实体类型ClusterSharding.start在集群的每个节点启动时典型做法用ClusterSharding.start注册支持的实体类型// 来源ClusterShardingSpec.scala #counter-start val counterRegion: ActorRef ClusterSharding(system).start( typeName Counter, entityProps Props[Counter](), settings ClusterShardingSettings(system), extractEntityId extractEntityId, extractShardId extractShardId)Java 版本// 来源ClusterShardingTest.java #counter-start OptionString roleOption Option.none(); ClusterShardingSettings settings ClusterShardingSettings.create(system); ActorRef startedCounterRegion ClusterSharding.get(system) .start(Counter, Props.create(Counter.class), settings, messageExtractor);注意ClusterSharding.start返回的引用就是该实体类型的ShardRegion可将其传给需要发送消息的代码。特别地当当前节点的角色与ClusterShardingSettings中指定的角色不匹配时ClusterSharding.start会以 Proxy Only 模式仅代理模式启动一个ShardRegion详见下文“Proxy Only Mode”。3. 消息提取器extractEntityId 与 extractShardIdextractEntityId与extractShardId是应用自定义的两个函数Java 中是MessageExtractor接口的entityId/entityMessage/shardId方法用于从入站消息中提取实体标识符与分片标识符// 来源ClusterShardingSpec.scala #counter-extractor val extractEntityId: ShardRegion.ExtractEntityId { case EntityEnvelope(id, payload) (id.toString, payload) case msg Get(id) (id.toString, msg) } val numberOfShards 100 val extractShardId: ShardRegion.ExtractShardId { case EntityEnvelope(id, _) (id % numberOfShards).toString case Get(id) (id % numberOfShards).toString case ShardRegion.StartEntity(id) // StartEntity 由 remembering entities 特性使用 (id.toLong % numberOfShards).toString case _ throw new IllegalArgumentException() }Java 版本ShardRegion.MessageExtractor见 ClusterShardingTest.java 的#counter-extractor片段。示例展示了两种在消息中定义实体标识符的方式Get消息自带标识符counterId/id字段EntityEnvelope包装标识符信封持有id真正发给实体的消息是payload字段。注意extractEntityId返回的二元组中第二部分或 Java 中entityMessage的返回值才是真正发送给实体 Actor 的消息——这使得在需要时可以解开信封。4. 分片Shard与分片算法设计一个 Shard 是一组会被统一管理的实体分组方式由extractShardId定义。对同一个实体标识符分片标识符必须始终保持一致否则实体 Actor 可能被意外地在多处同时启动。设计分片算法本身就是一个有趣的挑战目标是产生均匀分布每个分片中的实体数量大致相同。经验法则分片数量应约为计划最大集群节点数的十倍不必精确分片数少于节点数时部分节点将不承载任何分片分片数过多会导致分片管理效率下降例如再平衡rebalancing开销增大且由于 Coordinator 参与每个分片首条消息的路由延迟也会增加分片算法必须在运行中的集群所有节点上保持一致只能在停止全部节点后更改。大多数场景下够用的简单算法是实体标识符hashCode的绝对值对分片数取模。这个算法已由ShardRegion.HashCodeMessageExtractor提供无需自行实现。5. 通过 ShardRegion 发送消息发往实体的消息总是经由本地的ShardRegion// 来源ClusterShardingSpec.scala #counter-usage val counterRegion: ActorRef ClusterSharding(system).shardRegion(Counter) counterRegion ! Get(123) expectMsg(0) counterRegion ! EntityEnvelope(123, Increment) counterRegion ! Get(123) expectMsg(1)ClusterSharding.shardRegion可按实体类型名取回ShardRegion引用。ShardRegion在不知道分片位置时会查询 Coordinator 定位分片把消息委托到正确的节点并在实体的第一条消息到达时按需创建实体 Actoron demand。工作原理ShardRegion、Shard 与 ShardCoordinator完整的消息路由与故障处理场景描述见 cluster-sharding-concepts.mdTyped 文档目录下经典与 Typed 共享该概念文档其核心机制如下ShardRegion在每个集群节点或具有特定角色的节点组上启动持有两个应用自定义函数提取实体 ID 与分片 ID。首次收到某分片的消息时向中央协调者ShardCoordinator请求该分片的位置。Shard一组统一管理的实体。ShardRegion确认分片归属后创建Shard监督者作为子 Actor实体再由Shard按需创建。ShardCoordinator决定哪个ShardRegion拥有某个Shard以**集群单例Cluster Singleton**形式运行在集群最老成员或特定角色节点组的最老成员上。分片分配决策由此集中做出保证所有节点对分片位置有一致视图——这是集群中至多只有一个特定实体 Actor 实例的关键保障。消息路由的两种典型场景SCCoordinatorSRRegionSShardEEntity场景 1消息指向本 Region 的未知分片。M1到达SR1→ 映射到S1SR1向SC询问位置 →SC回答归属SR1→SR1创建子 ActorS1并转发 →S1创建E1并转发。此后SR1处理S1的消息无需再问SC。场景 2消息指向远端 Region 的分片。M2到达SR1→ 映射到S2询问SC→SC回答归属SR2→SR1将S2的缓冲消息发给SR2此后SR1直接转发S2消息到SR2。分片位置解析期间该分片的入站消息会被缓冲待分片归属确定后投递已解析分片的后续消息则直接送达目标不再经过 Coordinator。分片再平衡Rebalancing为使新加入节点被利用Coordinator 会发起分片迁移。过程为通知所有ShardRegion某分片开始 handoff此后该分片消息被缓冲Coordinator 也不应答该分片的位置请求→ 原 Region 向该分片内所有实体发送stopMessage默认PoisonPill停止实体 → 全部实体终止后原 Region 向 Coordinator 确认 handoff 完成 → Coordinator 应答新位置请求为分片分配新家缓冲消息投递到新位置。实体的状态不会随迁移被转移重要状态应通过 Persistence 持久化以便在新位置恢复。故障恢复Coordinator 节点崩溃或不可达并被 down 掉后新的 Coordinator 单例接管并恢复状态。故障期间已知位置的分片仍可用而新未知分片的消息会缓冲至新 Coordinator 就绪。消息顺序与投递语义只要发送方通过同一个ShardRegion向实体投递消息顺序即被保持在缓冲上限内消息以尽力而为的方式投递与普通消息发送一样是at-most-once语义。需要 at-least-once 端到端可靠投递时可叠加 Reliable Delivery 特性其#sharding章节专门说明与分片的配合。开销提示发往新分片或久未使用的分片的消息会因 Coordinator 往返引入额外延迟再平衡也会增加延迟——这应在设计分片解析时考虑避免分片过细。一旦分片位置已知唯一开销只是经由ShardRegion转发而非直发。Sharding 状态存储模式State Store ModeCluster Sharding 管理两类状态ShardCoordinator 状态——各Shard的位置强制存在Remembering Entities 状态——每个Shard中的实体列表可选默认关闭。针对这两类状态当前有两种存储模式模式底层技术说明Distributed Data 模式Akka Distributed DataCRDT默认模式Persistence 模式Akka Persistence事件溯源已废弃deprecated官方对 Persistence 模式的废弃说明见 includes/cluster.md 的#sharding-persistence-mode-deprecated片段Persistence 状态存储模式已废弃建议将 Coordinator 状态迁移到ddata若使用 remembering entities则迁移到eventsourced的 remember entities 存储。旧persistence模式写入的 remembered entities 数据可被新的eventsourcedremember entities 模式读取但一旦迁移就无法回退。更改状态存储模式需要 全集群重启不能滚动升级。Distributed Data 模式默认akka.cluster.sharding.state-store-mode ddataCoordinator 状态在集群内复制但非持久不落盘由 Distributed Data 以WriteMajority/ReadMajority一致性处理集群所有节点停止后状态不再需要并被丢弃集群分片在每个节点使用独立的 Distributed DataReplicator。若分片配置了角色则每个角色一个 Replicator从而支持某些实体类型只用节点子集、另一些用另一子集每个 Replicator 的名字包含节点角色因此角色配置必须在所有节点上一致滚动升级期间不能修改角色修改角色同样需要全集群重启akka.cluster.sharding.distributed-data配置段用于设置 Distributed Data 参数且不同实体类型不能使用不同的distributed-data设置。Persistence 模式已废弃akka.cluster.sharding.state-store-mode persistence由于运行在集群环境中Persistence 必须配置分布式日志distributed journal。该模式仍保留历史兼容用途新项目不应再使用。Proxy Only 模式仅代理ShardRegion也可以仅以代理模式启动自身不承载任何实体但知道如何把消息委托到正确位置。两种触发方式通过ClusterSharding.startProxy方法显式启动代理当ClusterSharding.start传入的ClusterShardingSettings中指定的角色与当前节点角色不匹配时start会自动以代理模式启动ShardRegion。典型用途是前端节点例如在 ClusterShardingSpec.scala 的多节点测试中sixth节点配置为frontend角色其余为backend它通过代理向远端的Counter实体发送消息测试断言消息实际由远端处理lastSender.path.address不等于本地地址。Java 中还可通过startProxy配合数据中心的Optional参数跨数据中心代理见#proxy-dc片段。Passivation实体钝化若实体状态已持久化可以停止不使用的实体以降低内存消耗。实现方式由应用决定例如在实体中定义接收超时context.setReceiveTimeout。一个重要的坑如果实体自行停止时邮箱中还有已入队的消息这些消息会被丢弃。要优雅钝化且不丢失消息实体应向父级Shard发送ShardRegion.PassivatePassivate中包装的停止消息会被发回给实体实体收到后应自行停止从收到Passivate到实体终止期间Shard会缓冲入站消息这些缓冲消息随后会被投递给实体的新化身new incarnation。在计数器示例中ReceiveTimeout120 秒触发context.parent ! Passivate(stopMessage Stop)实体收到Stop后调用context.stop(self)。Typed API 还提供了**自动钝化Automatic Passivation**策略体系详见 typed/cluster-sharding.md 的#automatic-passivation章节默认策略是空闲实体钝化2 分钟无消息也可配置基于活动实体上限的策略如 LRU / MRU / LFU / SLRU / TinyLFU 复合策略通过akka.cluster.sharding.passivation.strategy切换或strategy none关闭开启 Remembering Entities 时自动钝化会被禁用。相关默认配置见 reference.conf 的passivation段。Remembering Entities记住实体开启后分片在再平衡或崩溃恢复后会自动重建之前运行的实体未开启时实体只在消息到达时被重启。注意实体自身的状态不会被恢复除非实体本身已持久化例如使用事件溯源。开启方式// ClusterShardingSettings 中设置 rememberEntities true // 并要求 extractShardId 能处理 Shard.StartEntity(EntityId) // 即必须能从 EntityId 中提取出 ShardId// 来源ClusterShardingSpec.scala #extractShardId-StartEntity val extractShardId: ShardRegion.ExtractShardId { case EntityEnvelope(id, _) (id % numberOfShards).toString case Get(id) (id % numberOfShards).toString case ShardRegion.StartEntity(id) // StartEntity 由 remembering entities 特性使用 (id.toLong % numberOfShards).toString case _ throw new IllegalArgumentException() }Remember Entities 存储remember-entities-storeakka.cluster.sharding.remember-entities-store有两个取值ddata默认使用 Distributed Data 持久化到磁盘LMDB以支持全集群重启后实体被记住若不需要可设akka.cluster.sharding.distributed-data.durable.keys []关闭落盘例如 Kubernetes 无持久卷环境、或无需全集群重启后恢复。注意remember-entities-storeddata在Java 17下运行 LMDB 需要添加 JVM 参数--add-opensjava.base/sun.nio.chALL-UNNAMED --add-opensjava.base/java.nioALL-UNNAMED。eventsourced使用事件溯源存储活动分片与每个分片的活动实体必须配置 persistence 与 snapshot 插件akka.cluster.sharding.remember-entities-store eventsourced akka.cluster.sharding.journal-plugin-id plugin akka.cluster.sharding.snapshot-plugin-id plugin从废弃的 persistence 模式迁移未使用 remembering entities全集群重启后即可迁移到ddata使用 remembering entities 的两种迁移路径状态存储与 remember entities 都用ddata全集群重启后所有被记住的实体将丢失状态存储用ddata、remember entities 用eventsourced新的eventsourced存储能读取旧persistence模式写入的数据全集群重启后实体仍被记住。迁移现有 remembered entities 时还需在application.conf中为所用日志配置事件适配器以 Cassandra 日志为例akka.persistence.cassandra.journal { event-adapters { coordinator-migration akka.cluster.sharding.OldCoordinatorStateMigrationEventAdapter } event-adapter-bindings { akka.cluster.sharding.ShardCoordinator$Internal$DomainEvent coordinator-migration } }一旦迁移成功就不能回退到旧persistence存储因此也不支持滚动升级回滚。另外注意使用 Distributed Data 模式时实体标识符存储在 Distributed Data 的 Durable Storage 中。你可能需要修改akka.cluster.sharding.distributed-data.durable.lmdb.dir的配置因为默认目录包含 ActorSystem 的远程端口若使用动态端口0每次端口不同会导致此前存储的数据无法被加载。若不需要全集群重启后启动相同实体可关闭 durable 存储以换取更好性能akka.cluster.sharding.distributed-data.durable.keys []Supervision监督实体 Actor 默认使用重启restarting监督策略。如果需要为实体 Actor 使用其他supervisorStrategy必须创建一个中间父 Actor由它定义对子实体的监督策略// 来源ClusterShardingSpec.scala #supervisor class CounterSupervisor extends Actor { val counter context.actorOf(Props[Counter](), theCounter) override val supervisorStrategy OneForOneStrategy() { case _: IllegalArgumentException SupervisorStrategy.Resume case _: ActorInitializationException SupervisorStrategy.Stop case _: DeathPactException SupervisorStrategy.Stop case _: Exception SupervisorStrategy.Restart } def receive { case msg counter.forward(msg) } }然后以与启动实体 Actor 相同的方式启动该监督者// 来源ClusterShardingSpec.scala #counter-supervisor-start ClusterSharding(system).start( typeName SupervisedCounter, entityProps Props[CounterSupervisor](), settings ClusterShardingSettings(system), extractEntityId extractEntityId, extractShardId extractShardId)注意两点被停止的实体在有新消息指向它时会再次启动若使用 on stop 型退避监督策略Backoff supervisor必须设置并使用一个最终终止消息用于钝化详见 fault-tolerance.md 的#sharding小节。Graceful Shutdown优雅停机向ShardRegion发送ShardRegion.GracefulShutdownJavaShardRegion.gracefulShutdownInstance即可让该ShardRegion交接hand off其承载的所有分片然后ShardRegionActor 被停止可通过watch该ShardRegion获知交接完成。交接期间其他 Region 会像 Coordinator 触发再平衡时一样缓冲这些分片的消息分片停止后Coordinator 会在别处重新分配它们。该过程由 Coordinated Shutdown 自动执行因此它是集群成员优雅离开graceful leaving流程的一部分。移除内部 Cluster Sharding 数据此操作仅与 Persistence 模式相关Coordinator 存储了分片位置这些数据在重启整个 Akka 集群时即可安全移除注意这不包含应用数据。仓库提供独立工具类akka.cluster.sharding.RemoveInternalClusterShardingData以独立 Java 主程序方式运行java -classpath jar files, including akka-cluster-sharding akka.cluster.sharding.RemoveInternalClusterShardingData -2.3 entityType1 entityType2 entityType3程序包含在akka-cluster-shardingjar 中最简单的运行方式是使用与普通应用相同的 classpath 和配置也可从 sbt 或 Maven 类似方式运行程序参数为实体类型名与ClusterSharding.start中的typeName一致若第一个参数指定-2.3还会尝试移除 Akka 2.3.x 以不同 persistenceId 存储的数据。严格警告运行该程序时必须停止所有正在使用 Cluster Sharding 的集群节点绝不可在运行中的集群上执行典型场景是因错误的 downing provider 在网络分区时意外同时运行了两个集群导致 Coordinator 数据损坏无法启动。巡检 Cluster Sharding 状态提供两个状态查询请求用途是测试与监控而非提供直接向单个实体发消息的通道ShardRegion.GetShardRegionStateJavagetShardRegionStateInstance→ 返回ShardRegion.CurrentShardRegionStateJavaShardRegion.ShardRegionState包含某个 Region 中运行的分片标识符以及每个分片中存活的实体。ShardRegion.GetClusterShardingStats→ 查询集群中所有 Region返回ShardRegion.ClusterShardingStats包含每个 Region 中运行的分片标识符及各分片中存活实体的数量。若某些分片查询失败例如分片过忙未能在akka.cluster.sharding.shard-region-query-timeout内应答两个响应还会附带按 Region 归类的失败分片标识符集合。此外可通过ClusterSharding.shardTypeNamesJavagetShardTypeNames获取所有已启动分片的类型名。Lease租约—— 防止分片双跑的最后防线Lease 可作为一种额外的安全措施确保一个分片不会同时运行在两个节点上。出现双跑的可能原因包括网络分区且没有合适的 downing provider部署失误导致出现两个相互独立的 Akka 集群网络分区两侧移除成员与关闭成员之间的时序问题。使用方式设置akka.cluster.sharding.use-lease指向所用 lease 的配置位置。每个分片会尝试获取名为actor system name-shard-type name-shard id的 leaseowner 为Cluster(system).selfAddress.hostPort。若分片无法获取 lease它将保持未初始化状态属于它的实体消息会在ShardRegion中缓冲若初始化后 lease 丢失该Shard会被终止。配置总览ClusterShardingSettings是ClusterSharding.start方法的参数即每个实体类型都可以按需配置不同的设置。核心配置项位于 reference.conf经典与 reference.confTyped含number-of-shards要点如下配置项默认值说明akka.cluster.sharding.role实体运行在具有特定角色的节点上留空表示所有节点akka.cluster.sharding.remember-entitiesoff是否记住实体分片重启后自动重建akka.cluster.sharding.remember-entities-storeddataremember entities 存储ddata或eventsourcedakka.cluster.sharding.state-store-modeddataCoordinator 状态存储ddata默认或persistence已废弃akka.cluster.sharding.number-of-shards1000默认HashCodeMessageExtractor使用的分片数全集群必须一致join 时有配置检查修改需停止所有节点akka.cluster.sharding.passivation.strategydefault-idle-strategy自动钝化策略none/off关闭akka.cluster.sharding.coordinator-failure-backoff5 sCoordinator 状态存储失败后的重启退避指数退避至 5 倍akka.cluster.sharding.retry-interval2 sRegion 未获应答时重试注册与位置请求的间隔akka.cluster.sharding.buffer-size100000Region 缓冲消息的最大数量akka.cluster.sharding.handoff-timeout60 s分片再平衡交接超时akka.cluster.sharding.shard-start-timeout10 s等待 Region 确认承载分片的时间akka.cluster.sharding.entity-restart-backoff10 sremembering entities 下实体未用 Passivate 自停后的重启退避akka.cluster.sharding.rebalance-interval10 s周期性再平衡检查间隔akka.cluster.sharding.least-shard-allocation-strategy.rebalance-absolute-limit0单轮再平衡的最大分片数设为 0如 20可启用 2.6.10 引入的新再平衡算法推荐未来将成为默认akka.cluster.sharding.least-shard-allocation-strategy.rebalance-relative-limit0.1单轮再平衡分片数占已知分片总数的比例与 absolute-limit 取较小值akka.cluster.sharding.shard-region-query-timeout3 s查询 Region 全部分片状态的超时akka.cluster.sharding.entity-recovery-strategyall分片恢复实体策略all或constant固定速率akka.cluster.sharding.use-lease分片使用的 lease 配置路径空为不使用akka.cluster.sharding.verbose-debug-loggingoff逐消息级别调试日志生产环境慎开akka.cluster.sharding.coordinator-state.write-majority-plus/read-majority-plus3/5状态写入/读取要求的多数派之外的额外节点数可增容错但会变慢再平衡与节点增减默认的LeastShardAllocationStrategy将新分片分配给已分配分片最少的ShardRegion节点新节点加入时现有节点上的分片会向新节点再平衡同样从分片最多的 Region 选取、分配给最少的 Region。通过rebalance-absolute-limit/rebalance-relative-limit可限制每轮再平衡数量避免一次性迁移过多分片导致大量事件溯源实体同时启动的额外负载。最小成员数启动官方推荐配合akka.cluster.min-nr-of-members或akka.cluster.role.role-name.min-nr-of-members使用——分片分配会延迟到至少该数量的 Region 已启动并向 Coordinator 注册避免大量分片先被分配给首个注册节点、之后再迁移到其他节点。健康检查Cluster Sharding 内置与 Akka Management 兼容的健康检查注册后即返回 healthy默认开启可通过akka.management.health-checks.readiness-checks { sharding }关闭。每个 Region 的监控默认关闭需要时在akka.cluster.sharding.healthcheck.names中列出实体类型名如[counter-1, HelloWorld]健康检查在成员 up 后持续失败一段时间后会被禁用始终返回成功以免阻塞 Kubernetes 滚动升级。总结Cluster Sharding 是 Akka 在集群中水平扩展有状态 Actor 的标准方案。其核心要素可归纳为实体Entity 分片Shard Region/Coordinator 两级路由 可插拔的分片分配与状态存储。实践中需要重点把握提取器三件事extractEntityId、extractShardId必须对同一实体 ID 稳定返回同一分片 ID且全集群算法一致分片数量按节点数 × 10的经验法则设计number-of-shards修改需全集群停机状态持久化实体状态用事件溯源恢复Coordinator 状态默认走ddatapersistence模式已废弃生命周期治理用Passivate优雅钝化、GracefulShutdown优雅停机、状态查询消息做监控必要时用 Lease 兜底防双跑。进一步的阅读路径Typed API 完整文档 typed/cluster-sharding.md、概念详解 typed/cluster-sharding-concepts.md以及仓库中的多节点测试 ClusterShardingSpec.scala 与配置默认值 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),仅供参考
返回列表