Akka 集群实战:Actor 分布式、集群分片与容错

深入 Akka 与 Akka Cluster:Actor 系统与 ActorRef、监督与容错、集群成员管理与故障检测、集群分片(Cluster Sharding)、集群单例、分布式数据,以及如何用 Akka 构建可扩展的分布式 Actor 系统。

引言

Actor 模型解决了单机并发,但要支撑「水平扩展、高可用、跨机器分布」,需要把 Actor 系统升级为集群。Akka Cluster 正是为此而生:它让 Actor 可以透明地分布在多台机器上,一个 Actor 系统跨节点协作,故障节点自动摘除,消息路由自动重定向。

本文从单机 Akka 出发,一路走到生产级 Akka Cluster:先搭起 Actor 系统与 ActorRef,理解监督与容错(失败如何隔离与恢复);再进入集群,讲成员管理与故障检测(Split Brain 怎么处理)、集群分片(把海量 Actor 均匀分布到节点)、集群单例与分布式数据;最后给出一个「聊天室 + 用户会话分片」的综合实战。

前置:Actor 模型基础(https://plumephp.com/actor-model-detailed-explanation/)与函数式思维(https://plumephp.com/scala-functional-programming/)。如需 Scala 3 语法对照见 https://plumephp.com/scala3-modern-features/。


目录


1. Actor 系统与 ActorRef:单机 Akka 起步

1.1 依赖与最小系统

// build.sbt
libraryDependencies += "com.typesafe.akka" %% "akka-actor-typed" % "2.8.5"
import akka.actor.typed.{ActorSystem, Behavior, ActorRef}
import akka.actor.typed.scaladsl.Behaviors

object Greeter {
  sealed trait Command
  final case class Greet(name: String, replyTo: ActorRef[String]) extends Command

  def apply(): Behavior[Command] = Behaviors.receiveMessage { case Greet(name, replyTo) =>
    replyTo ! s"你好, $name"
    Behaviors.same
  }
}

val system: ActorSystem[Greeter.Command] =
  ActorSystem(Greeter(), "demo-system")

1.2 ActorRef 的意义

  • 拿到 ActorRef 就能 !(发消息),不关心 Actor 在哪台机器。
  • 这是「位置透明」的基础:集群正是靠这一抽象把 Actor 搬到多节点。

1.3 发送与接收

val probe: ActorRef[String] = ...   // 测试探针
val greeter: ActorRef[Greeter.Command] = system ? Greeter()
greeter ! Greeter.Greet("Alice", probe)

2. 监督与容错:失败隔离与恢复

2.1 监督树

ActorSystem
└── /user/greeter    (父)
    └── /user/greeter/child   (子)

父 Actor 监督子 Actor 的失败,决定重启/停止/升级。

2.2 监督策略

import akka.actor.typed.scaladsl.Behaviors
import akka.actor.typed.{SupervisorStrategy, Behavior}

def child(): Behavior[String] = Behaviors.receiveMessage {
  case "boom" => throw new RuntimeException("子崩溃")
  case msg    => println(msg); Behaviors.same
}

// 父:子失败后重启,最多 3 次
def parent(): Behavior[String] =
  Behaviors.supervise(child())
    .onFailure[RuntimeException](SupervisorStrategy.restart.withLimit(3, 1.second))

2.3 容错哲学

原则含义
失败隔离一个 Actor 崩溃不影响兄弟 Actor
重启恢复崩溃后按策略重启,状态可重建
策略上升本级处理不了,交父 Actor 处理

3. 集群成员管理与故障检测

3.1 集群是什么

多个独立 Actor 系统(节点)通过 gossip 协议组成一个逻辑集群。每个节点知道集群成员状态。

3.2 成员状态机

Joining → Up → Leaving → Exiting → Removed
                └── Down (故障)
  • Joining:刚加入,等待确认。
  • Up:正常运行。
  • Leaving:优雅退出。
  • Down:被判定故障。

3.3 故障检测:Phi 累加探测器

Akka 用 Phi Accrual Failure Detector 估计「节点存活的概率」,超阈值判定 Down。比固定心跳更鲁棒(容忍瞬时抖动)。

3.4 Split Brain 问题

网络分区时,两个子集群各自认为自己「拥有全部」,若都启动集群单例会冲突。解法:Split Brain Resolver(SBR)规则。

akka.cluster.split-brain-resolver.active-strategy = keep-majority

常用策略:keep-majority(保留多数派)、keep-oldest(保留最老节点)、static-quorum(静态法定人数)。


4. 集群分片:海量 Actor 的分布之道

4.1 问题

用户会话、订单等「有身份的海量 Actor」若每个都建一个进程级 Actor,机器再多也不够。**集群分片(Cluster Sharding)**把这些 Actor 按 id 分片,每个分片落到一个节点,路由自动寻址。

4.2 概念

Entity(实体)= 一个具体 Actor(如 UserSession)
Shard   = 一组 Entity 的容器,落在某节点
ShardRegion = 客户端入口,路由到正确分片

4.3 定义与使用

import akka.cluster.sharding.typed.scaladsl.{ClusterSharding, Entity, EntityTypeKey}
import akka.cluster.sharding.typed.ClusterShardingSettings

object UserSession {
  val TypeKey: EntityTypeKey[Command] = EntityTypeKey("UserSession")
  sealed trait Command
  final case class AddMessage(userId: Long, msg: String) extends Command

  def apply(userId: Long): Behavior[Command] = Behaviors.receiveMessage {
    case AddMessage(_, msg) => println(s"用户$userId 收到: $msg"); Behaviors.same
  }
}

val sharding = ClusterSharding(system)
val region = sharding.init(Entity(UserSession.TypeKey)(
  createBehavior = ctx => UserSession(ctx.entityId.toLong)))

// 客户端发送:自动路由到持有该实体的节点
region ! UserSession.AddMessage(userId = 42, msg = "hello")

4.4 分片数量与分布

  • 默认每个节点若干分片,分片按哈希分布。
  • 实体空闲可 passivate(休眠),节点数变化时自动 rebalance。

5. 集群单例与分布式数据

5.1 集群单例:全集群唯一

需要「全局只有一个」的组件(如总协调器、全局 ID 生成器)用 Cluster Singleton:

import akka.cluster.singleton.ClusterSingleton
import akka.cluster.singleton.ClusterSingletonSettings

val singleton = ClusterSingleton(system).init(
  ClusterSingletonSettings(system),
  Behaviors.receiveMessage[String] { case msg => println(s"单例: $msg"); Behaviors.same })

5.2 分布式数据(Distributed Data)

基于 CRDT 的键值存储,支持无协调冲突合并:

import akka.cluster.ddata.typed.scaladsl.{DistributedData, Replicator}
import akka.cluster.ddata.{LWWRegister, LWWRegisterKey}

val replicator = DistributedData(system).replicator
val key = LWWRegisterKey[Int]("counter")
replicator ! Replicator.Update(key, LWWRegister(0), Replicator.WriteLocal)(reg => reg.withValue(reg.value + 1))

5.3 选择

需求用
全局唯一组件Cluster Singleton
多节点共享低写频率数据Distributed Data
高吞吐实体路由Cluster Sharding

6. 消息投递语义与分布式陷阱

6.1 三种投递

语义说明适用
at-most-once可能丢失默认、可容忍
at-least-once不丢但可能重需幂等处理
exactly-once不丢不重成本高,少用

6.2 分布式陷阱清单

陷阱说明
消息顺序不同 Actor 间无全局顺序
延迟与超时需设置 ask 超时
序列化跨节点消息必须可序列化(Proto/JSON)
脑裂用 SBR 策略

6.3 序列化配置

akka.actor.serialization-bindings {
  "myapp.messages.**" = jackson-json
}

7. 综合实战:分布式会话聊天室

7.1 设计

客户端 → ShardRegion(按 userId 分片) → UserSession Actor
                                    └─ 订阅 Room Actor(集群单例或分片)

7.2 会话 Actor

object RoomActor {
  val TypeKey = EntityTypeKey[Command]("Room")
  sealed trait Command
  final case class Broadcast(from: Long, text: String) extends Command

  def apply(): Behavior[Command] = Behaviors.setup { ctx =>
    Behaviors.receiveMessage { case Broadcast(from, text) =>
      ctx.log.info(s"[${ctx.self.path}] $from: $text")
      Behaviors.same
    }
  }
}

val roomRegion = ClusterSharding(system).init(
  Entity(RoomActor.TypeKey)(_ => RoomActor()))

7.3 会话 Actor 转发到房间

def session(userId: Long, roomRef: ActorRef[RoomActor.Command]): Behavior[UserSession.Command] =
  Behaviors.receiveMessage {
    case UserSession.AddMessage(_, msg) =>
      roomRef ! RoomActor.Broadcast(userId, msg)
      Behaviors.same
  }

7.4 运行效果

多节点部署后,任意节点发消息,都路由到持有 Room/UserSession 的节点——位置透明性让业务代码完全不知道跨机器。


8. 生产化:配置、监控与运维

8.1 最少配置

akka {
  actor.provider = "cluster"
  remote.artery.canonical.port = 2551
  cluster {
    seed-nodes = ["akka://sys@node1:2551", "akka://sys@node2:2551"]
    roles = ["app"]
  }
}

8.2 监控

import akka.management.cluster.bootstrap.ClusterBootstrap
// 用 Akka Management 暴露集群状态端点

8.3 部署要点

项建议
种子节点3 个(奇数,利于脑裂决策)
健康检查就绪/存活探针
日志集中采集(结构化)
版本升级滚动升级,观察成员状态

9. 总结:何时用 Akka Cluster

9.1 适合的场景

  • 海量有身份实体(会话、订单、游戏角色)→ 分片。
  • 需要失败隔离与自动恢复 → Actor 监督。
  • 需要水平扩展与跨节点协作 → 集群。

9.2 不适合的场景

  • 简单 CRUD 服务 → 用传统框架更省。
  • 强事务需求 → Actor 不是数据库。
  • 团队不熟 Actor 模型 → 学习成本高。

9.3 一句话心法

Akka Cluster 的价值在「位置透明 + 失败隔离」:业务代码写的是单个 Actor,运行在整张集群网格上。


延伸阅读

  • https://plumephp.com/actor-model-detailed-explanation/ — Actor 模型理论基石
  • https://plumephp.com/scala-functional-programming/ — 不可变与消息传递的契合
  • [[erlang]] 专题 — Erlang/OTP 的 Actor 对照
  • Akka 官方文档

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「scala」更多文章

  1. 纯函数式效果系统实战:Cats Effect IO 与 ZIO
  2. Scala.js 与 Scala Native:跨平台编译、互操作与工程实践
  3. Scala 领域建模实战:ADT、类型驱动设计与模块化架构