十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Akka Persistence 插件机制完全指南:可插拔的 Journal、快照存储与持久化查询后端

Akka Persistence 插件机制完全指南:可插拔的 Journal、快照存储与持久化查询后端 后端并发编程异步编程【免费下载链接】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 是 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 storedurable 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.leveldbLevelDB 插件还需要额外声明如下依赖org.fusesource.leveldbjni % leveldbjni-all % 1.8LevelDB 文件的默认存储位置是当前工作目录下名为journal的目录可通过配置修改路径支持相对或绝对见 reference.conf 中的默认值dir journalakka.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://example127.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.confakka.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 云原生、使用 DynamoDBDynamoDB 插件注意其不支持 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 中定义的插件配置骨架classplugin-dispatcher编写自定义插件配置。深入理解 reference.conf 与 Persistence.scala 中的解析逻辑是掌握 Akka Persistence 插件机制的关键一步。赞分享后端并发编程异步编程【免费下载链接】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点击查看免费下载相关推荐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),仅供参考
返回列表