聊分布式系统时,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-streams | Reactive Streams 实现,带背压 |
akka-persistence | 事件溯源、Actor 状态持久化 |
akka-http | 高性能 HTTP 服务 |
演进脉络上需要知道三件事:
- Typed Actors(
Behavior<T>,2015 年引入)成为 Akka 2.6 的推荐 API,Classic API 进入维护模式——Typed 用类型系统约束了消息协议。 - 2022 年 Akka 2.7 改为 BSL 商业许可,不再纯开源。
- 社区因此分叉出 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
│
┌───────────────┼───────────────┐
▼ ▼ ▼
数据仓库 搜索索引 风控/通知
五个运行时环节:
- 入口写入:支付网关回调一个薄服务,只做校验然后把
PaymentConfirmed(orderId=42)写进 Kafka 就返回。写入端完全无状态,随便扩容。 - 消费路由:Akka 集群里跑一个 Kafka consumer group,消息按
orderId分区到达,每条消息转成 command 发给OrderEntity(42)。 - 实体处理:
OrderEntity(42)常驻内存,持有这笔订单的完整状态机(Created → Paid → Shipped),收到PaymentConfirmed后校验前置状态、产生领域事件OrderPaid。 - 事件持久化:事件追加写 journal(事件溯源)。Actor 崩溃后由事件重放恢复状态。
- 事件外发: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 几乎一一对应,文档和教程可以互相套用。