- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
Akka Persistence 是 Akka 提供的持久化扩展,其核心设计之一便是存储后端完全可插拔:事件日志(Journal)、快照存储(Snapshot Store)、持久化状态存储(Durable State Store)以及持久化查询(Persistence Query)都可以通过配置无缝替换为不同的实现。本文以官方文档 persistence-plugins.md 为主线,结合本仓库源码(akka-persistence模块)与测试代码,系统讲解官方维护的插件生态、各插件的功能边界、插件启用与预打包插件的完整配置与使用方式,帮助你在实际项目中正确选型并落地配置。
插件生态总览:官方维护的存储后端
在 Akka Persistence 扩展中,存储后端是可以插拔的,Akka 团队官方维护了以下四类持久化插件:
| 插件 | 适用数据库 | 特点 |
|---|---|---|
| R2DBC 插件 | PostgreSQL、H2(内存/文件模式)、Yugabyte 等响应式关系型数据库 | 支持 Akka Persistence 的最新特性,官方推荐优先于 JDBC 插件使用 |
| Cassandra 插件 | Apache Cassandra | 契合 Cassandra 数据模型,但不支持部分较新的持久化特性(详见下文) |
| AWS DynamoDB 插件 | AWS DynamoDB | 不支持 Durable State,也不支持“仅从最后一条事件恢复” |
| JDBC 插件 | 任何具备 JDBC 驱动的传统关系型数据库 | 适合已有 JDBC 技术栈的存量项目,新项目官方建议改用 R2DBC |
说明:上述插件的独立文档由各自项目维护,本文聚焦于 Akka 核心仓库内置的预打包插件(LevelDB Journal、本地快照存储、插件代理)以及通用的启用与初始化机制。
功能边界:Cassandra 与 JDBC 插件的特性限制
官方文档明确指出,Cassandra 插件与 JDBC 插件不支持 Akka Persistence 后续新增的部分特性,具体包括:
eventsBySlices查询- 基于 gRPC 的 Projections(投影)
- 基于 gRPC 的 Replicated Event Sourcing(复制式事件溯源)
- Projection 实例数量的动态伸缩
- 低延迟 Projections
- 从快照启动的 Projections
- 大量 Projections 的可扩展性
- Durable State 实体(JDBC 插件仅部分支持)
- 仅从最后一条事件恢复(Recovery from only last event)
与之对应,R2DBC 插件支持上述绝大多数最新特性,这也是官方在“新项目优先选择 R2DBC”这一建议背后的技术原因。DynamoDB 插件除了不支持 Durable State 之外,也不支持“仅从最后一条事件恢复”。在选型时,请先对照这张能力清单确认你的业务是否依赖其中某一项特性。
启用插件:默认配置与按 Actor 单独指定
插件有两种启用方式:为所有持久化 Actor 设置“默认”插件,或者由单个持久化 Actor 自己指定一套插件。
当持久化 Actor 没有覆写journalPluginId和snapshotPluginId方法时,持久化扩展会使用reference.conf中配置的“默认” journal、snapshot-store 与 durable-state 插件。在 akka-persistence 的 reference.conf 中,这些默认值都是空字符串,必须由你在用户侧的application.conf中显式覆写:
akka.persistence.journal.plugin = "" akka.persistence.snapshot-store.plugin = "" akka.persistence.state.plugin = ""从源码看,Persistence.scala 中defaultJournalPluginId与defaultSnapshotPluginId均为lazy val,且defaultSnapshotPluginId在未配置时会打印告警并回退到akka.persistence.no-snapshot-store(即NoSnapshotStore),而 journal 未配置则直接抛异常——这印证了文档中的提示:如果不使用快照,可以完全不配置快照存储插件。同时注意 reference.conf 中有一条重要提示:Cluster Sharding 内部使用快照,因此使用 Cluster Sharding 时必须配置快照存储插件。
- 将事件写入本地 LevelDB 的 journal 插件示例见后文 Local LevelDB journal;
- 将快照以独立文件写入本地文件系统的快照存储示例见后文 Local snapshot store;
- durable state store 相对较新,一个可用的实现是 Akka Persistence JDBC 插件。
在 PersistencePluginDocSpec.scala 的测试配置中还可以看到自定义插件的完整结构——插件配置项下必须声明class(插件实现类的全限定名,需提供无参构造器或仅接收一个com.typesafe.config.Config参数的构造器)与plugin-dispatcher(插件 Actor 使用的调度器):
akka.persistence.journal.plugin = "my-journal" # My custom journal plugin my-journal { # Class name of the plugin. class = "docs.persistence.MyJournal" # Dispatcher for the plugin actor. plugin-dispatcher = "akka.actor.default-dispatcher" }插件的急切初始化(Eager Initialization)
默认情况下,持久化插件是按需懒启动的——只有在被实际使用时才创建。但在某些场景下(例如希望提前完成插件 Actor 的启动、预热连接池),你可能希望插件在 ActorSystem 启动时立即初始化。做法分两步:
- 将
akka.persistence.Persistence加入akka.extensions键; - 在
akka.persistence.journal.auto-start-journals与akka.persistence.snapshot-store.auto-start-snapshot-stores下列出要自动启动的插件 ID。
例如,要为 leveldb journal 插件和本地快照存储插件做急切初始化:
akka { extensions = [akka.persistence.Persistence] persistence { journal { plugin = "akka.persistence.journal.leveldb" auto-start-journals = ["akka.persistence.journal.leveldb"] } snapshot-store { plugin = "akka.persistence.snapshot-store.local" auto-start-snapshot-stores = ["akka.persistence.snapshot-store.local"] } } }该机制的底层实现位于 Persistence.scala:扩展初始化时读取journal.auto-start-journals与snapshot-store.auto-start-snapshot-stores两个字符串列表,并依次调用journalFor(id)/snapshotStoreFor(id)强制实例化对应插件,同时输出Auto-starting journal plugin ...的日志。reference.conf中这两个键的默认值为空列表[],即默认不自动启动任何插件。
预打包插件:随 Akka Persistence 内置的四种实现
Akka Persistence 模块内置了少量持久化插件,但官方明确警告:这些插件均不适合在 Akka Cluster 中用于生产环境(原因在于它们都依赖本地文件系统或单点共享)。它们主要服务于单机开发、测试与教学场景。
Local LevelDB journal
该插件将事件写入本地 LevelDB 实例。
⚠️ 警告:LevelDB 插件不能用于 Akka Cluster,因为其存储位于本地文件系统中。
LevelDB journal 已弃用(deprecated),官方不建议用其构建新应用,推荐用 Akka Persistence JDBC 作为替代。插件配置入口为akka.persistence.journal.leveldb,通过如下配置启用(见 PersistencePluginDocSpec.scala):
# Path to the journal plugin to be used akka.persistence.journal.plugin = "akka.persistence.journal.leveldb"LevelDB 插件还需要额外声明如下依赖:
"org.fusesource.leveldbjni" % "leveldbjni-all" % "1.8"LevelDB 文件的默认存储位置是当前工作目录下名为journal的目录,可通过配置修改,路径支持相对或绝对(见 reference.conf 中的默认值dir = "journal"):
akka.persistence.journal.leveldb.dir = "target/journal"使用该插件时,每个 ActorSystem 都会运行自己独立的 LevelDB 实例。
关于删除与压缩(Compaction)的特殊性:LevelDB 有一个显著特点——删除操作并不会真正从 journal 中移除消息,而是为每条被删除的消息追加一条“墓碑记录(tombstone)”。在高频删除的重度使用场景下,journal 文件会持续膨胀。为此,LevelDB 提供了专门的 journal 压缩功能,通过以下配置按 persistence id 设定触发阈值(完整示例):
# Number of deleted messages per persistence id that will trigger journal compaction akka.persistence.journal.leveldb.compaction-intervals { persistence-id-1 = 100 persistence-id-2 = 200 # ... persistence-id-N = 1000 # use wildcards to match unspecified persistence ids, if any "*" = 250 }compaction-intervals的默认值为空(reference.conf),即默认不压缩;你可以为每个 persistence id 单独设置阈值,也可以用"*"通配符兜底匹配所有未显式指定的 persistence id。此外 reference.conf 还暴露了fsync = on(写入时是否 fsync)、checksum = off(读取时是否校验 checksum)、native = on(使用 JNI 原生 LevelDB 还是 Java 移植版)等可调参数。
Shared LevelDB journal
共享 LevelDB journal 同样已弃用并将从未来的 Akka 版本中移除,不建议新应用使用。在多节点环境中做测试时,官方推荐使用inmemjournal 配合 Persistence Plugin Proxy;当然,生产环境实际使用的插件同样是好的测试选择。
说明:该插件已被 Persistence Plugin Proxy 取代。
共享 LevelDB 实例通过实例化SharedLeveldbStoreActor 启动。Scala 版本(测试代码):
import akka.persistence.journal.leveldb.SharedLeveldbStore val store = system.actorOf(Props[SharedLeveldbStore](), "store")Java 版本(测试代码):
final ActorRef store = system.actorOf(Props.create(SharedLeveldbStore.class), "store");默认情况下,共享实例将 journaled 消息写入当前工作目录下名为journal的本地目录,可通过配置修改存储位置(示例):
akka.persistence.journal.leveldb-shared.store.dir = "target/shared"使用共享 LevelDB 存储的 ActorSystem 必须激活akka.persistence.journal.leveldb-shared插件(示例):
akka.persistence.journal.plugin = "akka.persistence.journal.leveldb-shared"该插件必须通过注入(远程)SharedLeveldbStoreActor 引用来完成初始化,注入方式是调用SharedLeveldbJournal.setStore方法并传入 Actor 引用。Scala 的典型用法是先通过actorSelection定位 store 并用Identify握手,收到ActorIdentity后再注入(完整片段):
import akka.actor._ trait SharedStoreUsage extends Actor { override def preStart(): Unit = { context.actorSelection("akka://example@127.0.0.1:2552/user/store") ! Identify(1) } def receive = { case ActorIdentity(1, Some(store)) => SharedLeveldbJournal.setStore(store, context.system) } }Java 版本采用同样的Identify/ActorIdentity握手模式(LambdaPersistencePluginDocTest.java)。内部 journal 命令(由持久化 Actor 发送)会一直缓冲,直到注入完成;注入是幂等的,即只有第一次注入生效。
Local snapshot store
该插件将快照文件写入本地文件系统。
⚠️ 警告:本地快照存储插件不能用于 Akka Cluster,因为其存储位于本地文件系统中。
插件配置入口为akka.persistence.snapshot-store.local,通过如下配置启用(示例):
# Path to the snapshot store plugin to be used akka.persistence.snapshot-store.plugin = "akka.persistence.snapshot-store.local"默认存储位置是当前工作目录下名为snapshots的目录(reference.conf 中dir = "snapshots"),可通过配置修改,路径支持相对或绝对:
akka.persistence.snapshot-store.local.dir = "target/snapshots"再次强调:配置快照存储插件并不是强制性的。如果你不使用快照,就无需配置它。此外 reference.conf 中还提供了两个值得留意的进阶参数:max-load-attempts = 3(最新快照恢复失败时,依次回退尝试更旧的快照文件的最大次数)与snapshot-is-optional = false(若为true,快照加载失败时忽略快照、回放全部事件来恢复,但切勿在删除了事件的情况下开启,否则会导致恢复出的状态错误)。
Persistence Plugin Proxy
用于测试目的的持久化插件代理,允许在同一节点(或不同节点)上的多个 ActorSystem 之间共享同一个 journal 与快照存储。例如,它可以让持久化 Actor 故障转移到备用节点,并从备用节点继续使用共享的 journal 实例。其工作原理是:将所有的 journal/snapshot store 消息转发给同一个共享的持久化插件实例,因此它支持被代理插件所支持的任何用例。
⚠️ 警告:共享的 journal/snapshot store 是单点故障,只应出于测试目的使用。
journal 与 snapshot store 代理分别通过akka.persistence.journal.proxy与akka.persistence.snapshot-store.proxy配置项控制。完整配置骨架见 reference.conf:
akka.persistence.journal.proxy { class = "akka.persistence.journal.PersistencePluginProxy" plugin-dispatcher = "akka.actor.default-dispatcher" # 在托管目标 journal 的 ActorSystem 配置中设为 on start-target-journal = off # 目标 journal 的插件配置路径 target-journal-plugin = "" # 其他节点连接代理时使用的地址(可选) target-journal-address = "" # 目标查找的初始化超时 init-timeout = 10s } akka.persistence.snapshot-store.proxy { class = "akka.persistence.journal.PersistencePluginProxy" plugin-dispatcher = "akka.actor.default-dispatcher" start-target-snapshot-store = off target-snapshot-store-plugin = "" target-snapshot-store-address = "" init-timeout = 10s }使用步骤可以归纳为三点:
- 指定目标插件:将
target-journal-plugin或target-snapshot-store-plugin键设置为要使用的底层插件(例如akka.persistence.journal.inmem); - 选定托管节点:在恰好一个ActorSystem 中,将
start-target-journal与start-target-snapshot-store设为on——该系统将实例化共享的持久化插件; - 告知代理目标位置:通过配置键
target-journal-address/target-snapshot-store-address,或编程方式调用PersistencePluginProxy.setTargetLocation方法。
编程方式的底层实现(PersistencePluginProxy.scala)会向 journal 与快照存储代理发送TargetLocation(address)消息;PersistencePluginProxy.start则通过journalFor(null)/snapshotStoreFor(null)强制实例化代理。
注意:Akka 会懒启动扩展,代理也不例外。为了让代理正常工作,目标节点上的持久化插件必须被实例化。可以通过实例化
PersistencePluginProxyExtension扩展(参见 extending-akka.md),或调用PersistencePluginProxy.start方法来实现。
注意:被代理的持久化插件可以(也应该)使用其原本的配置键进行配置。
小结与选型建议
| 场景 | 推荐方案 |
|---|---|
新项目、需要最新特性(eventsBySlices、gRPC Projections、复制式事件溯源、Durable State 等) | R2DBC 插件(PostgreSQL / H2 / Yugabyte) |
| 已有 Cassandra 基础设施、不依赖最新特性 | Cassandra 插件 |
| AWS 云原生、使用 DynamoDB | DynamoDB 插件(注意其不支持 Durable State 与“仅恢复最后事件”) |
| 存量 JDBC 技术栈迁移 | JDBC 插件(新项目不推荐) |
| 单机开发 / 测试 / 教学 | 内置 LevelDB journal、本地快照存储、inmemjournal |
| 多节点共享 journal 的测试 | Persistence Plugin Proxy(仅测试用途,单点故障) |
无论选择哪种后端,启用机制都是统一的:通过akka.persistence.journal.plugin、akka.persistence.snapshot-store.plugin与akka.persistence.state.plugin三个键完成默认绑定,必要时配合auto-start-journals/auto-start-snapshot-stores实现插件急切初始化,并遵循 reference.conf 中定义的插件配置骨架(class+plugin-dispatcher)编写自定义插件配置。深入理解 reference.conf 与 Persistence.scala 中的解析逻辑,是掌握 Akka Persistence 插件机制的关键一步。
- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
相关推荐
Go Micro Model 数据模型层完全指南:结构化 CRUD、查询与可插拔存储后端
Go Micro Model 数据模型层完全指南:结构化 CRUD、查询与可插拔存储后端 导读 model 包是 Go Micro( go micro.dev/
后端微服务AI AgentRPC框架为 Akka Persistence Durable State 构建存储后端插件:从接口实现到配置激活的完整指南
为 Akka Persistence Durable State 构建存储后端插件:从接口实现到配置激活的完整指南 导读 本文面向希望为 Akka 的持久化扩展
后端并发编程异步编程Akka Persistence Query 完全指南:基于 Durable State 的 CQRS 查询端实战
Akka Persistence Query 完全指南:基于 Durable State 的 CQRS 查询端实战 导读 本文聚焦 Akka 中面向 Durab
后端并发编程异步编程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考