Akka Persistence 存储后端插件开发指南:Journal 与 Snapshot Store 的构建、配置与 TCK 验证
【免费下载链接】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
导读
Akka Persistence 扩展中的Journal(事件日志)与Snapshot Store(快照存储)都是可插拔的存储后端。本文以官方文档 persistence-journals.md 为核心,系统讲解如何从零构建一个新的存储后端:包括实现AsyncWriteJournal/SnapshotStore插件 API、通过 HOCON 配置激活插件、理解插件 Actor 的执行模型与构造器约定,以及借助官方 TCK(Technology Compatibility Kit)规范验证插件的正确性与性能。读完本文,你将掌握编写一个生产级 Akka Persistence 存储插件所需的全部 API 细节、配置要点与测试方法论。
一、可插拔的存储后端架构
Akka Persistence 将"事件日志存储"与"快照存储"抽象为两个独立的后端接口,应用可以:
- 实现插件 API提供自己的存储实现(如 MySQL、Cassandra、文件系统、内存等);
- 通过配置激活任意插件,无需修改业务代码。
官方文档同时指出,Akka 社区项目页维护了一份现成的 journal 与 snapshot store 插件目录,社区插件列表见 Community plugins(外部链接,仅作背景参考)。本文关注的是如何自行构建存储后端。
开发一个插件,首先需要引入以下导入(官方文档要求的插件开发导入):
Scala(示例见 PersistencePluginDocSpec.scala):
import akka.persistence._ import akka.persistence.journal._ import akka.persistence.snapshot._Java(示例见 LambdaPersistencePluginDocTest.java):
import akka.dispatch.Futures; import akka.persistence.*; import akka.persistence.journal.japi.*; import akka.persistence.snapshot.japi.*;从源码结构看,Scala 插件 API 位于akka.persistence.journal与akka.persistence.snapshot包(见 AsyncWriteJournal.scala),Java 插件 API 则位于akka.persistence.journal.japi与akka.persistence.snapshot.japi包(见 AsyncWritePlugin.java)。
二、Journal 插件 API:继承 AsyncWriteJournal
一个 journal 插件需要继承AsyncWriteJournal。从源码看,AsyncWriteJournal是一个同时混入了WriteJournalBase与AsyncRecovery的Actor(AsyncWriteJournal.scala),它已经实现了消息协议处理(receiveWriteJournal)、电路断路器(circuit breaker)包装、消息重排序(Resequencer)等框架逻辑,插件作者只需要实现下面三个核心方法。
2.1 需要实现的方法
AsyncWriteJournal要求插件实现的方法(源码注释见 AsyncWriteJournal.scala):
| 方法 | 职责 | 保护机制 |
|---|---|---|
asyncWriteMessages(messages: immutable.Seq[AtomicWrite]): Future[immutable.Seq[Try[Unit]]] | 异步批量写入持久化消息 | 受 circuit breaker 保护 |
asyncDeleteMessagesTo(persistenceId: String, toSequenceNr: Long): Future[Unit] | 异步删除指定persistenceId上截至toSequenceNr(含)的全部消息 | 受 circuit breaker 保护 |
receivePluginInternal: Actor.Receive | 可选覆盖,用于处理自定义内部消息(如f pipeTo self实现高级功能) | — |
Java API 对应实现(AsyncWritePlugin.java):
Future<Iterable<Optional<Exception>>> doAsyncWriteMessages(Iterable<AtomicWrite> messages); Future<Void> doAsyncDeleteMessagesTo(String persistenceId, long toSequenceNr);2.2 asyncWriteMessages 的语义约定
这是 journal 插件最重要、也最容易出错的接口,源码注释给出了明确的契约:
- 批量仅为性能优化:
Seq中的消息不要求原子写入,底层存储支持批量插入(batch insert)时通常能获得更高吞吐; - AtomicWrite 的原子性:每个
AtomicWrite要么包含单条PersistentRepr(来自persist),要么包含多条PersistentRepr(来自persistAll),其中的所有PersistentRepr必须原子写入,即全有或全无。若底层存储不支持多条事件的原子写,应返回Try的Failure(携带UnsupportedOperationException说明问题),且应在插件文档中声明此限制; - 成功/失败的界定:Future 只有在批量中所有消息都被确认持久化(后续回放可见)时才以成功完成;任何存储失败、或对"是否已写入"存在不确定时,都必须以失败完成;数据存储连接问题必须以 Future 失败信号,而不是消息拒绝(rejection)信号;
- 按 persistenceId 保持顺序:对同一
persistenceId的实际写入应当串行化,否则可能出现"后写的事件先对查询端/回放可见"的不一致;PersistentActor在前一次WriteMessages完成前不会发出新的写请求; - sender 已置空:
PersistentRepr中的sender字段已被置为ActorRef.noSender,避免为回放时已失效的引用浪费存储空间; - 最高序号并发读取:对同一
persistenceId的最高序号查询可能与写操作并发执行(例如重启的 Actor 在其未完成写入前就开始恢复),此时应尽量推迟读取最高序号直到未完成写入全部结束,否则PersistentActor可能复用序号。
2.3 同步阻塞型后端的写法
如果存储后端只支持同步、阻塞式写入,官方文档建议把阻塞调用包在Future中实现。Scala 示例(PersistencePluginDocSpec.scala):
class MyJournal extends AsyncWriteJournal { def asyncWriteMessages(messages: immutable.Seq[AtomicWrite]): Future[immutable.Seq[Try[Unit]]] = Future.fromTry(Try { // blocking call here ??? }) def asyncDeleteMessagesTo(persistenceId: String, toSequenceNr: Long): Future[Unit] = ??? def asyncReplayMessages(persistenceId: String, fromSequenceNr: Long, toSequenceNr: Long, max: Long)( replayCallback: (PersistentRepr) => Unit): Future[Unit] = ??? def asyncReadHighestSequenceNr(persistenceId: String, fromSequenceNr: Long): Future[Long] = ??? // optionally override: override def receivePluginInternal: Receive = super.receivePluginInternal }Java 示例(LambdaPersistencePluginDocTest.java):
class MyAsyncJournal extends AsyncWriteJournal { @Override public Future<Iterable<Optional<Exception>>> doAsyncWriteMessages( Iterable<AtomicWrite> messages) { try { Iterable<Optional<Exception>> result = new ArrayList<Optional<Exception>>(); // blocking call here... // result.add(..) return Futures.successful(result); } catch (Exception e) { return Futures.failed(e); } } @Override public Future<Void> doAsyncDeleteMessagesTo(String persistenceId, long toSequenceNr) { return null; } @Override public Future<Void> doAsyncReplayMessages( String persistenceId, long fromSequenceNr, long toSequenceNr, long max, Consumer<PersistentRepr> replayCallback) { return null; } @Override public Future<Long> doAsyncReadHighestSequenceNr(String persistenceId, long fromSequenceNr) { return null; } }三、恢复接口:AsyncRecovery
journal 插件还必须实现AsyncRecovery中定义的两个方法,用于消息回放与序号恢复。该 trait 的完整契约见 AsyncRecovery.scala:
| 方法 | 语义要点 |
|---|---|
asyncReplayMessages(persistenceId, fromSequenceNr, toSequenceNr, max)(recoveryCallback) | 异步回放消息,逐条调用replayCallback;Future 在全部匹配消息回放完成后完成,任一回放失败则以失败完成。已被标记删除的消息也必须传给replayCallback(此时deleted返回true)。fromSequenceNr = -1表示只回放最后一条消息(如果有)。该调用不受 circuit breaker 保护(可能耗时很长),插件必须自行防护后端存储无响应,并在合理时间内完成 Future——不允许忽略完成 |
asyncReadHighestSequenceNr(persistenceId, fromSequenceNr) | 异步读取指定persistenceId的最高存储序号;PersistentActor在恢复后以此为起点写入新事件;该序号也是后续asyncReplayMessages的toSequenceNr上限(除非用户指定更低值)。journal 必须维护最高序号且永不降低它。受 circuit breaker 保护。fromSequenceNr是搜索起点提示:恢复时若使用了快照则为快照序号,否则为0L |
一个重要的框架行为(从 AsyncWriteJournal.scala 的receiveWriteJournal可以看出):框架默认先调用asyncReadHighestSequenceNr取得highSeqNr,再将toSequenceNr截断为min(toSequenceNr, highSeqNr)后调用asyncReplayMessages。
3.1 可选优化:AsyncReplay
源码还提供了AsyncReplaytrait 作为优化路径(AsyncRecovery.scala):若插件实现了replayMessages(将回放与读取最高序号合并为一次调用),则AsyncRecovery中的两个方法将不再被调用。这在需要避免"先读最高序号、再回放"两次后端往返的场景下可以显著提升恢复性能。
四、Journal 插件的最小配置与激活
官方文档给出的 journal 插件最小激活配置(PersistencePluginDocSpec.scala):
# Path to the journal plugin to be used 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" }配置要点:
akka.persistence.journal.plugin指定激活的插件配置路径;顶层my-journal就是该插件自己的配置段;class是插件类的全限定名(FQCN);plugin-dispatcher是插件 Actor 使用的 dispatcher,未指定时默认使用akka.actor.default-dispatcher。
4.1 构造器约定
插件类必须具有以下三种签名之一的构造器:
- 一个
com.typesafe.config.Config参数 + 一个String参数(后者接收插件配置路径); - 只有一个
com.typesafe.config.Config参数; - 无参构造器。
构造器第一个参数收到的将是 ActorSystem 配置中该插件段(plugin section)的配置,String参数则是该插件的配置路径。这与 reference.conf 中journal-plugin-fallback的注释一致:插件类必须有无参构造器或仅含一个com.typesafe.config.Config参数的构造器。
4.2 插件 Actor 的执行模型
- journal 插件实例本身是一个 Actor,因此来自 PersistentActor 的请求对应的方法按顺序串行执行(由 Actor 邮箱天然保证);
- 插件可以委托给异步库、派生 Future、或委派给其他 Actor 来获得并行度;
- 切勿在系统默认 dispatcher 上运行 journal 任务/Future,那可能饿死其他任务——应使用插件自己的 dispatcher 或专用线程池。
4.3 插件可继承的默认配置
所有 journal 插件配置都会落到akka.persistence.journal.journal-plugin-fallback的默认值上(reference.conf),插件实现者与使用者都应了解:
| 配置项 | 默认值 | 说明 |
|---|---|---|
class | "" | 插件类 FQCN,必填 |
plugin-dispatcher | "akka.actor.default-dispatcher" | 插件 Actor 的 dispatcher |
replay-dispatcher | "akka.actor.default-dispatcher" | 消息回放使用的 dispatcher |
max-message-batch-size | 200 | 已废弃(无实际作用),PersistentActor 会写入自上次写入以来累积的全部消息 |
recovery-event-timeout | 30s | 恢复过程中两个事件之间超过该时间则恢复失败;注意它同样影响读取快照后再回放事件的过程(虽然配置在 journal 段下) |
circuit-breaker.max-failures | 10 | 断路器打开前的连续失败次数 |
circuit-breaker.call-timeout | 10s | 单次调用超时 |
circuit-breaker.reset-timeout | 30s | 断路器打开后尝试重置的等待时间 |
replay-filter.mode | repair-by-discard-old | 回放过滤器模式(详见第七章) |
replay-filter.window-size | 100 | 回放分析用的前瞻缓冲区大小(事件数) |
replay-filter.max-old-writers | 10 | 记住的旧 writerUuid 数量 |
replay-filter.debug | off | 开启后对每个回放事件输出详细调试日志 |
注意:
AsyncWriteJournal的源码中会为asyncWriteMessages、asyncDeleteMessagesTo、asyncReadHighestSequenceNr包上 circuit breaker(AsyncWriteJournal.scala),而asyncReplayMessages明确不受其保护。
五、Snapshot Store 插件 API:继承 SnapshotStore
快照存储插件必须继承SnapshotStoreActor(SnapshotStore.scala),并实现四个方法:
| 方法 | 职责 | 语义要点 |
|---|---|---|
loadAsync(persistenceId, criteria): Future[Option[SelectedSnapshot]] | 异步加载快照 | 返回None表示无快照(此时会回放全部事件);快照加载失败必须以 Future 失败完成,因为事件可能已被删除,仅回放事件可能无法得到有效状态;受 circuit breaker 保护 |
saveAsync(metadata, snapshot): Future[Unit] | 异步保存快照 | 受 circuit breaker 保护;框架会为metadata填充System.currentTimeMillis作为时间戳 |
deleteAsync(metadata): Future[Unit] | 删除指定快照 | 受 circuit breaker 保护 |
deleteAsync(persistenceId, criteria): Future[Unit] | 按条件删除快照 | 受 circuit breaker 保护 |
receivePluginInternal | 可选覆盖 | 处理插件自定义内部消息 |
Java API 对应实现(SnapshotStorePlugin.java):
Future<Optional<SelectedSnapshot>> doLoadAsync(String persistenceId, SnapshotSelectionCriteria criteria); Future<Void> doSaveAsync(SnapshotMetadata metadata, Object snapshot); Future<Void> doDeleteAsync(SnapshotMetadata metadata); Future<Void> doDeleteAsync(String persistenceId, SnapshotSelectionCriteria criteria);框架行为佐证(SnapshotStore.scala):当criteria == SnapshotSelectionCriteria.None时直接返回空结果,不调用插件;保存失败时会尝试deleteAsync(metadata)清理失败的快照;保存/删除成功或失败事件会先交给receivePluginInternal,再转发给 PersistentActor。
快照存储插件的最小激活配置(PersistencePluginDocSpec.scala):
# Path to the snapshot store plugin to be used akka.persistence.snapshot-store.plugin = "my-snapshot-store" # My custom snapshot store plugin my-snapshot-store { # Class name of the plugin. class = "docs.persistence.MySnapshotStore" # Dispatcher for the plugin actor. plugin-dispatcher = "akka.actor.default-dispatcher" }Snapshot store 插件与 journal 插件遵守相同的构造器约定(Config+String / Config / 无参三种之一),同样将插件段配置传入构造器;plugin-dispatcher未指定时默认akka.actor.default-dispatcher。快照插件的默认配置来自snapshot-store-plugin-fallback(reference.conf),其 circuit breaker 默认值(max-failures = 5、call-timeout = 20s、reset-timeout = 60s)与 journal 不同,配置时需注意区分。
同样的执行模型警告同样适用:snapshot store 实例是 Actor,请求串行执行;可以委托异步库/Future/其他 Actor 获取并行度;不要在系统默认 dispatcher 上运行快照存储任务/Future。
六、插件 TCK:用官方兼容性测试套件验证插件
为了帮助开发者构建正确、高质量的存储插件,Akka 提供了TCK(Technology Compatibility Kit,技术兼容性套件)。TCK 同时适用于 Java 与 Scala 项目,使用时需要引入akka-persistence-tck依赖。
6.1 引入依赖
需要在项目中加入akka-persistence-tck依赖(并确保从 Akka 的安全库仓库以 token 化 URL 获取依赖)。构建脚本位于仓库的 akka-persistence-tck 模块,应用侧大致如下:
// sbt libraryDependencies += "com.typesafe.akka" %% "akka-persistence-tck" % AkkaVersion % Test<!-- Maven --> <dependency> <groupId>com.typesafe.akka</groupId> <artifactId>akka-persistence-tck_2.13</artifactId> <version>${akka.version}</version> <scope>test</scope> </dependency>(AkkaVersion对应本文档目标 Akka 版本,具体版本号以当前仓库project/Dependencies.scala中定义的版本为准。)
6.2 Journal TCK:JournalSpec / JavaJournalSpec
在测试套件中直接继承JournalSpec(Scala)或JavaJournalSpec(Java)即可将 journal TCK 测试纳入测试套件。Scala 示例(PersistencePluginDocSpec.scala):
import akka.persistence.journal.JournalSpec class MyJournalSpec extends JournalSpec( config = ConfigFactory.parseString("""akka.persistence.journal.plugin = "my.journal.plugin"""")) { override def supportsRejectingNonSerializableObjects: CapabilityFlag = false // or CapabilityFlag.off override def supportsSerialization: CapabilityFlag = true // or CapabilityFlag.on }Java 示例(LambdaPersistencePluginDocTest.java):
@RunWith(JUnitRunner.class) class MyJournalSpecTest extends JavaJournalSpec { public MyJournalSpecTest() { super(ConfigFactory.parseString( "akka.persistence.journal.plugin = \"akka.persistence.journal.leveldb-shared\"")); } @Override public CapabilityFlag supportsRejectingNonSerializableObjects() { return CapabilityFlag.off(); } }6.3 能力开关(CapabilityFlag)
TCK 中部分测试是可选的:通过覆盖supports...方法,向 TCK 声明插件支持哪些能力,TCK 据此决定运行哪些测试。这些方法可以使用布尔值或CapabilityFlag.on/CapabilityFlag.off实现。从 JournalSpec.scala 看,TCK 基类默认的能力包括:
| 能力开关 | 默认值 | 含义 |
|---|---|---|
supportsSerialization | true | 插件支持事件序列化 |
supportsMetadata | false | 插件支持持久化元数据 |
supportsReplayOnlyLast | false | 插件支持只回放最后一条消息(fromSequenceNr = -1场景) |
supportsRejectingNonSerializableObjects | 由插件覆盖 | 插件能拒绝不可序列化对象(写入阶段以 rejection 信号返回) |
supportsAtomicPersistAllOfSeveralEvents | true | 支持persistAll产生的多条事件原子写入(覆盖于JournalSpec方法) |
6.4 性能粗测:JournalPerfSpec
TCK 还提供简单的基准类JournalPerfSpec(Java 为JavaJournalPerfSpec,见 JavaJournalPerfSpec.scala)。它包含JournalSpec的全部测试,并在 journal 上执行一些耗时操作、打印性能统计。官方文档明确说明:它不旨在提供规范的基准测试环境,但可用于在最典型场景下对 journal 性能获得粗略感受。
6.5 Snapshot Store TCK:SnapshotStoreSpec
将快照存储 TCK 测试纳入套件需继承SnapshotStoreSpec(PersistencePluginDocSpec.scala):
import akka.persistence.snapshot.SnapshotStoreSpec class MySnapshotStoreSpec extends SnapshotStoreSpec( config = ConfigFactory.parseString(""" akka.persistence.snapshot-store.plugin = "my.snapshot-store.plugin" """)) { override def supportsSerialization: CapabilityFlag = true // or CapabilityFlag.on }Java 对应为JavaSnapshotStoreSpec(LambdaPersistencePluginDocTest.java)。
6.6 生命周期钩子:beforeAll / afterAll
如果插件需要一些初始化/清理工作(如启动 mock 数据库、删除临时文件),可以覆盖beforeAll与afterAll钩入测试生命周期。注意必须调用super。Scala 示例(PersistencePluginDocSpec.scala):
class MyJournalSpec extends JournalSpec(config = ConfigFactory.parseString(""" akka.persistence.journal.plugin = "my.journal.plugin" """)) { override def supportsRejectingNonSerializableObjects: CapabilityFlag = true // or CapabilityFlag.on val storageLocations = List( new File(system.settings.config.getString("akka.persistence.journal.leveldb.dir")), new File(config.getString("akka.persistence.snapshot-store.local.dir"))) override def beforeAll(): Unit = { super.beforeAll() storageLocations.foreach(FileUtils.deleteRecursively) } override def afterAll(): Unit = { storageLocations.foreach(FileUtils.deleteRecursively) super.afterAll() } }Java 示例见 LambdaPersistencePluginDocTest.java(使用system().settings().config()读取存储路径并在beforeAll/afterAll中递归清理)。
官方文档强烈建议将这些 TCK 规范纳入测试套件:它们覆盖了从零编写插件时容易遗漏的大量边界场景(如
persistAll原子写、删除后回放、persistAsync批量、恢复超时、旧版deleted标志等)。
七、事件日志损坏与回放过滤器
如果 journal 无法阻止用户用相同的persistenceId并发运行多个 PersistentActor,事件日志很可能被破坏——表现为存在相同序号的事件。
官方文档给出的建议是:journal 在恢复时仍应交付这些事件,以便由回放过滤器(replay-filter)以与具体 journal 无关(journal agnostic)的方式决定如何处理。
回放过滤器正是AsyncWriteJournal内置的机制:从源码(AsyncWriteJournal.scala)可以看到,框架会在恢复时读取replay-filter配置,按mode创建ReplayFilter包装回复目标:
akka.persistence.journal.<plugin>.replay-filter { # 检测到无效事件时的处理方式: # `repair-by-discard-old` : 丢弃来自旧 writer 的事件,并记录警告(默认值) # `fail` : 使回放失败,记录错误 # `warn` : 记录警告但原样发出事件 # `off` : 完全禁用该功能 mode = repair-by-discard-old # 使用前瞻缓冲区分析事件,此处定义缓冲区大小(事件数) window-size = 100 # 记住多少个旧 writerUuid max-old-writers = 10 # 开启后对每个回放事件输出详细调试日志 debug = off }(默认值与详细说明见 reference.conf。)
ReplayFilter通过检查事件序号(sequence numbers)与 writerUuid来检测损坏的事件流:当检测到同一序号来自多个 writer 时,按mode采取丢弃旧 writer 事件、使回放失败或仅告警等策略。这样,即使底层 journal 无法从物理层面杜绝并发写同一persistenceId导致的日志损坏,应用层仍能以统一、可配置的方式恢复或发现损坏。
八、总结:构建存储后端的完整步骤
结合官方文档与仓库源码,一个完整的 Akka Persistence 存储后端开发流程如下:
- 实现 Journal:继承
AsyncWriteJournal(Java 为AsyncWriteJournal+japi接口),实现asyncWriteMessages、asyncDeleteMessagesTo,以及AsyncRecovery的asyncReplayMessages、asyncReadHighestSequenceNr;同步后端用Future.fromTry/Futures.successful包装阻塞调用; - 实现 Snapshot Store:继承
SnapshotStore(Java 为SnapshotStore+japi接口),实现loadAsync、saveAsync、两个deleteAsync; - 配置激活:提供
class(满足 Config+String / Config / 无参三种构造器约定之一)与plugin-dispatcher,并通过akka.persistence.journal.plugin/akka.persistence.snapshot-store.plugin指向插件配置段; - 遵守执行模型:插件是 Actor、请求串行;不要占用系统默认 dispatcher;利用
receivePluginInternal处理自定义消息; - 接入 TCK:引入
akka-persistence-tck,继承JournalSpec/JavaJournalSpec、SnapshotStoreSpec/JavaSnapshotStoreSpec,用CapabilityFlag声明能力,用beforeAll/afterAll管理环境,可选接入JournalPerfSpec做性能粗测; - 处理损坏日志:确保恢复时仍交付同序号事件,交由
replay-filter(repair-by-discard-old/fail/warn/off)统一决策。
仓库中可供继续研读的参考实现包括:官方文档测试样例 PersistencePluginDocSpec.scala 与 LambdaPersistencePluginDocTest.java、TCK 基类 JournalSpec.scala、以及内置插件的默认配置段(reference.conf 中的inmem、leveldb、local快照存储等),它们共同构成了插件开发最完整的参考闭环。
【免费下载链接】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),仅供参考