← blog 技术原理与实践 · 2026-09-22

Akka 与 MQ:Actor 模型、消息队列与两者的组合架构

从 Actor 模型讲清 Akka 的核心设计(mailbox、supervision、sharding、事件溯源),辨析 Akka 与 MQ 在解耦层级、状态位置、交付语义上的本质区别,并以电商订单系统为例完整拆解两者组合使用的架构、关键模式与接缝难题。

15 min read

聊分布式系统时,Akka 和 MQ(Kafka、RabbitMQ)经常被放在一起比较——它们都靠”异步消息”工作,看起来思路同源。但一个是最早解决并发正确性的编程模型,一个是最早解决系统间解耦的中间件基建。这篇文章先把 Akka 本身讲透,再辨析两者边界,最后用一个电商订单系统的完整例子展示两者如何组合——因为实践中它们不是二选一,而是上下游分工。

一、Akka 是什么

Akka 是 JVM 上的开源工具包和运行时(Scala/Java),用于构建高并发、分布式、容错的应用。它的理论基础是 Actor 模型——与 OOP、CSP 并列的并发编程模型,最早由 Carl Hewitt 在 1973 年提出,在 Erlang 语言上被电信行业验证了几十年。

1.1 Actor 模型的核心规则

Actor 模型可以浓缩成几条公理:

  • 一切都是 Actor:每个 Actor 是一个轻量级并发单元,拥有私有状态和一个邮箱(mailbox)。
  • 不共享内存:Actor 之间唯一的通信方式是异步消息(tell / !),没有锁、没有共享可变变量。
  • 一次处理一条消息:Actor 从邮箱逐条取消息处理,因此内部状态天然无需同步——并发问题的根源被直接移除了。
  • 顺序保证:同一发送者发给同一接收者的消息保证有序;不同发送者之间的到达顺序不定。

Actor 因此可以理解为三合一的结构:一个单消费者队列 + 一个处理器 + 一份私有状态。Akka 的 actor 极其轻量——官方数据是 4GB 内存可以创建约 650 万个 actor,远比线程便宜。

1.2 容错:Let it crash

Akka 最有辨识度的设计来自电信行业的自愈理念:

与其用 try-catch 包住一切可能出错的代码,不如让故障的 Actor 直接崩溃,由它的 supervisor(通常是创建它的父 Actor)决定如何处置。

处置策略有四种:

策略含义
Restart重启子 Actor(最常用),清理状态,恢复到干净起点
Resume完全忽略,继续处理下一条消息
Stop永久停止该子 Actor
Escalate自己也处理不了,向上一级上抛

这形成了 supervision 层级:故障被限制在子树内,不会污染整个系统。一个 Actor 挂了,受影响的只是它负责的那一小块业务,父级按策略重启即可——“系统永不停止”就是这么来的。

1.3 位置透明与集群

Actor 之间通过 ActorRef 互发消息,而 ActorRef 不区分目标是本地 JVM 还是远程节点:

// 本地和远程的代码完全一样
orderRef ! PaymentConfirmed(orderId = 42)

这让”单机程序”可以几乎无改动地扩展成集群。Akka Cluster 在其上补齐了分布式基建:成员管理(gossip 协议)、故障检测、分片(Cluster Sharding)、分布式发布订阅。其中 Cluster Sharding 是实战中最重要的组件,后面组合架构会反复用到它。

1.4 主要模块与演进

模块作用
akka-actor核心:Classic / Typed Actors
akka-remote远程通信
akka-cluster集群成员管理、Sharding、PubSub
akka-streamsReactive Streams 实现,带背压
akka-persistence事件溯源、Actor 状态持久化
akka-http高性能 HTTP 服务

演进脉络上需要知道三件事:

  1. Typed Actors(Behavior<T>,2015 年引入)成为 Akka 2.6 的推荐 API,Classic API 进入维护模式——Typed 用类型系统约束了消息协议。
  2. 2022 年 Akka 2.7 改为 BSL 商业许可,不再纯开源。
  3. 社区因此分叉出 Apache Pekko——Akka 2.6 的 Apache 2.0 许可分支,API 几乎一一对应。现在追求纯开源的项目一般选 Pekko。

题外话:今天的 akka.io 已经变成了 Lightbend 更名后的”企业级 Agentic AI 平台”,与经典开源 Akka actor 库是两个东西,查文档时注意别混。

二、Akka 与 MQ 的区别:同源不同职

直觉上”两者都是异步消息”,但把差别拆开看,它们几乎在每个维度上都相反。

2.1 相同点

两者共享同一个抽象根源:用消息传递替代共享内存。每个 Actor 的 mailbox 就是一个单消费者微型队列;Kafka 的 partition 也可以看作一个”带持久化的超大 mailbox”。Erlang 进程、Actor、MQ 的 queue,是同一个思想在不同尺度的投影。

2.2 本质:编程模型 vs 中间件

  • Actor 模型是并发编程模型。它要解决的首要问题是进程内/集群内的并发正确性——消灭锁、消灭共享可变状态。解耦只是它的副产品。
  • MQ 是中间件基建。它解耦的是身份与时间:生产者不需要知道消费者是谁、有几个、什么时候上线(发布订阅),也不需要消费者此刻活着(削峰、暂存)。

一个容易忽略的细节:Actor 之间其实是拓扑耦合的——A 必须持有 B 的 ActorRef 才能发消息,谁给谁发在代码里写死。MQ 反而做到了真正的身份解耦。所以严格说:Actor 解耦的是”并发逻辑与业务逻辑”,MQ 解耦的是”发送方与接收方”,方向并不一样。

2.3 状态:一个在计算单元里,一个外置

这是两者最深刻的区别,直接决定了各自的扩展模型:

  • Actor 有状态:状态在 Actor 内存里,处理下一条消息不需要查库。代价是——崩溃重启丢状态(需要事件溯源恢复)、不能随便水平扩容(状态钉在实例上,要靠 Sharding 按 entity id 分片)。
  • MQ 消费者无状态:随时加机器、随便重启,水平扩展是平凡的。代价是状态外置到数据库,每条消息处理都要先去拉状态。

注意一个常见误解:“MQ 无状态”要加限定——MQ 对业务无状态,但 broker 对消息本身有状态(持久化、offset、堆积位置)。

2.4 全维度对比

Actor (Akka)MQ (Kafka / RabbitMQ)
本质并发编程模型中间件基建
作用域进程内 / 集群内跨系统
状态在 Actor 里(内存)在 DB / 消费者侧
扩展方式位置透明 + Sharding加消费者 / 分区,天然平凡
消息持久默认否,挂了即丢默认是
顺序粒度同发送者→同接收者有序队列 / 分区内有序
交付语义尽力而为,靠应用层补at-least-once 起步,配套重试 / 死信 / 回放
失败哲学let it crash + supervisor 重启broker 暂存 + 消费重试 + DLQ 兜底
解耦对象并发逻辑与业务逻辑发送方与接收方

一句话总结:MQ 是哑管道,Actor 是智能端点。MQ 保证消息可靠地流过边界,Actor 保证消息被有状态地正确处理。

三、组合架构:以电商订单系统为例

实践中两者不是竞争关系而是上下游分工。典型拓扑是 MQ 出现两次、Akka 在中间:

外部系统 ──► MQ(入口缓冲/削峰)──► Akka 集群(有状态处理)──► MQ(事件外发)──► 下游

下面用一个电商订单系统完整走一遍。

3.1 架构总览

                    Kafka topics (入口)                Akka 集群
支付系统 ──► payment-events ─┐                 ┌─────────────────────────┐
库存服务 ──► stock-events ────┼──► Consumer ──► │ Cluster Sharding         │
用户操作 ──► commands ────────┘     Group       │  OrderEntity(orderId)×N  │
                                          │     │  状态机+私有状态          │
                                          │     └───────────┬─────────────┘
                                          │                 │ 事件追加写
                                          │                 ▼
                                          │          Journal (Cassandra/PG)
                                          │                 │
                                          │         Akka Projections
                                          │                 │
                                          └──► Kafka topic: order-lifecycle
                                                          │
                                          ┌───────────────┼───────────────┐
                                          ▼               ▼               ▼
                                       数据仓库          搜索索引          风控/通知

五个运行时环节:

  1. 入口写入:支付网关回调一个薄服务,只做校验然后把 PaymentConfirmed(orderId=42) 写进 Kafka 就返回。写入端完全无状态,随便扩容。
  2. 消费路由:Akka 集群里跑一个 Kafka consumer group,消息按 orderId 分区到达,每条消息转成 command 发给 OrderEntity(42)。
  3. 实体处理:OrderEntity(42) 常驻内存,持有这笔订单的完整状态机(Created → Paid → Shipped),收到 PaymentConfirmed 后校验前置状态、产生领域事件 OrderPaid。
  4. 事件持久化:事件追加写 journal(事件溯源)。Actor 崩溃后由事件重放恢复状态。
  5. 事件外发:Akka Projections 从 journal 读新事件、管理偏移量,发布到 Kafka 的 order-lifecycle,下游各自独立消费。

3.2 Sharding:百万实体怎么住进一个集群

百万订单不可能一个 Actor 全处理,也不可能百万线程。Cluster Sharding 把实体按 hash 分布到集群节点:

OrderEntity("order-42")  ──hash(orderId)──►  node 3
OrderEntity("order-99")  ──hash(orderId)──►  node 1
  • 实体按需创建:第一次收到消息时才实例化;
  • 闲置钝化(passivation):超时没消息的实体从内存卸载,下次消息来了再从 journal 重建——内存占用只和活跃实体数相关;
  • 节点宕机:它承载的分片被其他节点接管,实体从 journal 重放恢复——“let it crash” 从进程级扩展到了机器级;
  • 调用方完全无感:consumer 只管 sharding.entityRefFor("Order", id).tell(cmd)。

3.3 接缝处的关键对偶:partition key ≈ sharding key

整个组合架构成立的前提,是一个漂亮的对称关系:

Kafka 的 partition key ≈ Akka 的 sharding entity id

  • Kafka 保证:同 key 的消息进同一分区,分区内有序;
  • Akka Sharding 保证:同 entity id 的消息路由到同一实体,实体内串行处理。

两边用同一个业务 id(orderId)拼起来,就实现了”跨集群、全局范围内、对同一业务实体的串行处理”——这正是订单/账务类系统最需要又最难自己实现的东西。反过来,如果上游各系统没有统一用业务 id 做 key,顺序就断了,而且极难排查。

3.4 Outbox 模式:入口侧的一致性

入口侧有个经典问题:服务要”写 DB + 发 Kafka”两个操作,中间崩了就一边有一边没有(Kafka 不参与分布式事务)。解法是 outbox:

业务事务内:  写业务表 + 写 outbox 表   ← 同一事务, 原子
单独 relay:  轮询 outbox ──► 发 Kafka ──► 标记已发

在 Akka 的事件溯源体系里,journal 本身就是 outbox——事件落 journal 是原子的,Projections 负责往外转发。这是事件溯源架构的一个隐藏福利:不需要再单独建 outbox 表。

3.5 Akka Projections:出口侧的 offset 管理

Projections 干的事和 Kafka consumer 的 offset 管理同构,方向相反(journal → 外部系统):

  • 每个 projection(“发布到 Kafka”、“更新读模型表”……)独立记录偏移量,各自重放互不影响;
  • at-least-once 交付:处理完消息才存 offset,崩溃恢复后可能重复,要求 handler 幂等;
  • 支持 atLeastOnce / grouped(攒批)/ flow(响应式流)三种处理形态。

3.6 用代码感受一下(Pekko Typed 风格)

// 1. Kafka consumer 分片路由——partition key 和 sharding key 是同一个 orderId
val sharding = ClusterSharding(system)
Consumer.plainSource(settings, subscription)
  .map { record =>
    val orderId = record.key()            // Kafka 分区键
    val cmd     = decode(record.value())
    sharding.entityRefFor(OrderEntity.TypeKey, orderId) ! cmd
  }
  .run()

// 2. 有状态实体:状态机 + 事件溯源
object OrderEntity {
  def apply(orderId: String): Behavior[Command] =
    EventSourcedBehavior(
      persistenceId = PersistenceId("Order", orderId),
      emptyState    = OrderState.empty,
      commandHandler = (state, cmd) => cmd match {
        case ConfirmPayment(replyTo) if state.status == Created =>
          Effect.persist(OrderPaid(orderId)).thenReply(replyTo)(_ => Accepted)
        case ConfirmPayment(replyTo) =>
          Effect.reply(replyTo)(Rejected(s"当前状态 ${state.status} 不允许支付"))
      },
      eventHandler = (state, evt) => state.apply(evt)  // 重放时也走这里
    )
}

EventSourcedBehavior 把命令校验、事件产生、状态重建统一到一个结构里:commandHandler 决定”是否允许、持久化什么事件”,eventHandler 决定”状态如何随事件演进”。实体崩溃后框架自动回放事件重建状态,业务代码里看不到任何恢复逻辑。

3.7 接缝难题:决定成败的四个细节

1. 交付语义:at-least-once + 幂等

从 Kafka 到 Actor 是 at-least-once:consumer 处理了、还没提交 offset 就崩,重启后重放。实体必须幂等。事件溯源天然帮了大忙:处理 command 前查”这个事件是否已处理过”即可去重;没有事件溯源就要自己记处理过的 message key。

2. 顺序:两头拼接处最脆

同 orderId 的消息在 Kafka 分区内有序、在 Sharding 实体内串行——链路中段是稳的。脆的是两端:多个上游 topic 的 key 策略不统一、第三方系统发消息没用业务 id、或者 topic 扩分区时 key 重分布,都会悄悄断序。约定全链路 key 规范是架构问题,不是代码问题。

3. 背压:压不动 broker

Akka Streams 有完整的背压机制,能让 Actor 到 consumer 之间减速。但 Kafka broker 本身不会减速,partition 堆积照旧。所以背压的真正落地是运维策略:监控 consumer lag,lag 大了扩 consumer / 扩分片。

4. 重放:journal 是源,Kafka 是副本

这是组合架构里最深的坑。Kafka 的 log 可重放,事件溯源实体重放的却是 journal——两者不是一回事。原则:

journal 是唯一事实源,Kafka 是事实的分发副本。

新下游服务从头消费 Kafka 没问题(Kafka 保留全量);但要重建 Akka 实体状态只能靠 journal,不能”重新消费一遍 Kafka”(去重语义不同)。方向定成这样,一致性问题都有答案;反过来(以 Kafka 为源)会很痛苦。

3.8 决策清单

决策点建议
分区 key / sharding key全链路统一用同一个业务 id
幂等实体级去重 + 事件溯源,别指望 exactly-once
实体钝化设合理 passivation 超时,活跃常驻、冷实体卸载
事务边界入口侧 outbox,出口侧 Projections + offset
死信Akka 侧消费失败的消息单独进 DLQ topic,别静默重试
监控consumer lag + shard 分布 + 实体钝化率,三者都要
事件 schema用 schema registry(Avro/Protobuf)带演进规则

四、什么场景值得上这套

值得——实体的状态机复杂、数量多、且必须串行处理:

  • 订单 / 账务 / 交易:状态机迁移规则多,同单必须串行;
  • 物联网设备管理:每设备一实体,高频上报;
  • 游戏 / 协作文档:每房间 / 每文档一实体,实时性要求高;
  • 风控画像:每用户一实体,状态高频增量更新。

不值得——Actor 集群的运维复杂度(集群形成、分片再平衡、journal 运维)配不上收益:

  • 纯数据管道(ETL):Kafka Streams / Flink 更简单;
  • 简单 CRUD:MQ + 无状态服务 + DB 就够。

结语

回到最初的对比:Akka 与 MQ 的关系,不是”两种消息系统”,而是不同层级的两个答案——MQ 回答”消息如何可靠地跨过系统边界”,Actor 回答”消息如何在边界之内被有状态地正确处理”。理解了 partition key ≈ sharding key 这个接缝处的对偶关系,两者的组合就不再是两套技术的拼盘,而是一条连贯的、从边界到实体的设计线索。

最后一点选型提醒:这套架构的官方实现是 Akka(商业许可)的 akka-cluster-typed + akka-projection-kafka,纯开源等价物是 Apache Pekko + pekko-projection-kafka,API 几乎一一对应,文档和教程可以互相套用。

Sources

No external sources for this entry.

Related