ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Akka Persistence 存储后端插件开发指南:Journal 与 Snapshot Store 的构建、配置与 TCK 验证

Akka Persistence 存储后端插件开发指南:Journal 与 Snapshot Store 的构建、配置与 TCK 验证 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 的执行模型与构造器约定以及借助官方 TCKTechnology Compatibility Kit规范验证插件的正确性与性能。读完本文你将掌握编写一个生产级 Akka Persistence 存储插件所需的全部 API 细节、配置要点与测试方法论。一、可插拔的存储后端架构Akka Persistence 将事件日志存储与快照存储抽象为两个独立的后端接口应用可以实现插件 API提供自己的存储实现如 MySQL、Cassandra、文件系统、内存等通过配置激活任意插件无需修改业务代码。官方文档同时指出Akka 社区项目页维护了一份现成的 journal 与 snapshot store 插件目录社区插件列表见 Community plugins外部链接仅作背景参考。本文关注的是如何自行构建存储后端。开发一个插件首先需要引入以下导入官方文档要求的插件开发导入Scala示例见 PersistencePluginDocSpec.scalaimport akka.persistence._ import akka.persistence.journal._ import akka.persistence.snapshot._Java示例见 LambdaPersistencePluginDocTest.javaimport 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.scalaJava 插件 API 则位于akka.persistence.journal.japi与akka.persistence.snapshot.japi包见 AsyncWritePlugin.java。二、Journal 插件 API继承 AsyncWriteJournal一个 journal 插件需要继承AsyncWriteJournal。从源码看AsyncWriteJournal是一个同时混入了WriteJournalBase与AsyncRecovery的ActorAsyncWriteJournal.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.javaFutureIterableOptionalException doAsyncWriteMessages(IterableAtomicWrite messages); FutureVoid 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.scalaclass 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.javaclass MyAsyncJournal extends AsyncWriteJournal { Override public FutureIterableOptionalException doAsyncWriteMessages( IterableAtomicWrite messages) { try { IterableOptionalException result new ArrayListOptionalException(); // blocking call here... // result.add(..) return Futures.successful(result); } catch (Exception e) { return Futures.failed(e); } } Override public FutureVoid doAsyncDeleteMessagesTo(String persistenceId, long toSequenceNr) { return null; } Override public FutureVoid doAsyncReplayMessages( String persistenceId, long fromSequenceNr, long toSequenceNr, long max, ConsumerPersistentRepr replayCallback) { return null; } Override public FutureLong doAsyncReadHighestSequenceNr(String persistenceId, long fromSequenceNr) { return null; } }三、恢复接口AsyncRecoveryjournal 插件还必须实现AsyncRecovery中定义的两个方法用于消息回放与序号恢复。该 trait 的完整契约见 AsyncRecovery.scala方法语义要点asyncReplayMessages(persistenceId, fromSequenceNr, toSequenceNr, max)(recoveryCallback)异步回放消息逐条调用replayCallbackFuture 在全部匹配消息回放完成后完成任一回放失败则以失败完成。已被标记删除的消息也必须传给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是插件类的全限定名FQCNplugin-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-dispatcherakka.actor.default-dispatcher插件 Actor 的 dispatcherreplay-dispatcherakka.actor.default-dispatcher消息回放使用的 dispatchermax-message-batch-size200已废弃无实际作用PersistentActor 会写入自上次写入以来累积的全部消息recovery-event-timeout30s恢复过程中两个事件之间超过该时间则恢复失败注意它同样影响读取快照后再回放事件的过程虽然配置在 journal 段下circuit-breaker.max-failures10断路器打开前的连续失败次数circuit-breaker.call-timeout10s单次调用超时circuit-breaker.reset-timeout30s断路器打开后尝试重置的等待时间replay-filter.moderepair-by-discard-old回放过滤器模式详见第七章replay-filter.window-size100回放分析用的前瞻缓冲区大小事件数replay-filter.max-old-writers10记住的旧 writerUuid 数量replay-filter.debugoff开启后对每个回放事件输出详细调试日志注意AsyncWriteJournal的源码中会为asyncWriteMessages、asyncDeleteMessagesTo、asyncReadHighestSequenceNr包上 circuit breakerAsyncWriteJournal.scala而asyncReplayMessages明确不受其保护。五、Snapshot Store 插件 API继承 SnapshotStore快照存储插件必须继承SnapshotStoreActorSnapshotStore.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.javaFutureOptionalSelectedSnapshot doLoadAsync(String persistenceId, SnapshotSelectionCriteria criteria); FutureVoid doSaveAsync(SnapshotMetadata metadata, Object snapshot); FutureVoid doDeleteAsync(SnapshotMetadata metadata); FutureVoid 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 插件遵守相同的构造器约定ConfigString / Config / 无参三种之一同样将插件段配置传入构造器plugin-dispatcher未指定时默认akka.actor.default-dispatcher。快照插件的默认配置来自snapshot-store-plugin-fallbackreference.conf其 circuit breaker 默认值max-failures 5、call-timeout 20s、reset-timeout 60s与 journal 不同配置时需注意区分。同样的执行模型警告同样适用snapshot store 实例是 Actor请求串行执行可以委托异步库/Future/其他 Actor 获取并行度不要在系统默认 dispatcher 上运行快照存储任务/Future。六、插件 TCK用官方兼容性测试套件验证插件为了帮助开发者构建正确、高质量的存储插件Akka 提供了TCKTechnology 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 groupIdcom.typesafe.akka/groupId artifactIdakka-persistence-tck_2.13/artifactId version${akka.version}/version scopetest/scope /dependencyAkkaVersion对应本文档目标 Akka 版本具体版本号以当前仓库project/Dependencies.scala中定义的版本为准。6.2 Journal TCKJournalSpec / JavaJournalSpec在测试套件中直接继承JournalSpecScala或JavaJournalSpecJava即可将 journal TCK 测试纳入测试套件。Scala 示例PersistencePluginDocSpec.scalaimport 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.javaRunWith(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 能力开关CapabilityFlagTCK 中部分测试是可选的通过覆盖supports...方法向 TCK 声明插件支持哪些能力TCK 据此决定运行哪些测试。这些方法可以使用布尔值或CapabilityFlag.on/CapabilityFlag.off实现。从 JournalSpec.scala 看TCK 基类默认的能力包括能力开关默认值含义supportsSerializationtrue插件支持事件序列化supportsMetadatafalse插件支持持久化元数据supportsReplayOnlyLastfalse插件支持只回放最后一条消息fromSequenceNr -1场景supportsRejectingNonSerializableObjects由插件覆盖插件能拒绝不可序列化对象写入阶段以 rejection 信号返回supportsAtomicPersistAllOfSeveralEventstrue支持persistAll产生的多条事件原子写入覆盖于JournalSpec方法6.4 性能粗测JournalPerfSpecTCK 还提供简单的基准类JournalPerfSpecJava 为JavaJournalPerfSpec见 JavaJournalPerfSpec.scala。它包含JournalSpec的全部测试并在 journal 上执行一些耗时操作、打印性能统计。官方文档明确说明它不旨在提供规范的基准测试环境但可用于在最典型场景下对 journal 性能获得粗略感受。6.5 Snapshot Store TCKSnapshotStoreSpec将快照存储 TCK 测试纳入套件需继承SnapshotStoreSpecPersistencePluginDocSpec.scalaimport 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 对应为JavaSnapshotStoreSpecLambdaPersistencePluginDocTest.java。6.6 生命周期钩子beforeAll / afterAll如果插件需要一些初始化/清理工作如启动 mock 数据库、删除临时文件可以覆盖beforeAll与afterAll钩入测试生命周期。注意必须调用super。Scala 示例PersistencePluginDocSpec.scalaclass 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继承AsyncWriteJournalJava 为AsyncWriteJournaljapi接口实现asyncWriteMessages、asyncDeleteMessagesTo以及AsyncRecovery的asyncReplayMessages、asyncReadHighestSequenceNr同步后端用Future.fromTry/Futures.successful包装阻塞调用实现 Snapshot Store继承SnapshotStoreJava 为SnapshotStorejapi接口实现loadAsync、saveAsync、两个deleteAsync配置激活提供class满足 ConfigString / 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-filterrepair-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),仅供参考
返回列表