Akka Persistence Query 实战指南:用统一异步流接口构建 CQRS 读侧查询
Akka Persistence Query 实战指南用统一异步流接口构建 CQRS 读侧查询【免费下载链接】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-coreAkka Persistence Query 是 Akka 持久化体系的读侧查询接口它不关心事件如何写入 journal而是为各种 journal 插件定义了一套统一的、基于异步流Reactive Streams的查询协议让开发者可以用完全相同的 API 从 Cassandra、JDBC 数据库、LevelDB 等不同数据源订阅事件流。本文以 akka-docs/src/main/paradox/persistence-query.md 为主干结合仓库中的源码与测试PersistenceQuery.scala、Offset.scala、EventEnvelope.scala 等完整讲解依赖引入、ReadJournal 获取、预定义查询类型、物化值、读侧投影Materialized View、自定义查询插件开发以及集群横向扩展读完即可在自己的 Akka 应用中落地 CQRS 读侧。依赖引入Persistence Query 作为独立模块发布需要显式添加依赖。官方推荐通过 Akka BOM 统一管理版本// sbt libraryDependencies com.typesafe.akka %% akka-persistence-query % AkkaVersion!-- Maven -- dependency groupIdcom.typesafe.akka/groupId artifactIdakka-persistence-query_2.13/artifactId version${akka.version}/version /dependency// Gradle dependencies { implementation platform(com.typesafe.akka:akka-bom_2.13:${akka.version}) implementation com.typesafe.akka:akka-persistence-query_2.13 }引入该模块会自动带入 Akka Persistence 模块因为查询接口建立在写入 journal 的事件基础之上。注意本仓库中的 Akka 依赖需要通过 Akka 官方安全仓库的 tokenized URL 获取。引言Persistence Query 在 CQRS 中的定位Akka Persistence Query 是对 Event Sourcing事件溯源 的补充它提供一种通用的、基于异步流的查询接口由各种 journal 插件实现以暴露各自的查询能力。其最典型的应用场景是 CQRSCommand Query Responsibility Segregation命令查询职责分离架构中的查询侧又称读侧应用的写侧例如用 Akka Persistence 实现与查询侧完全分离。需要澄清的是Akka Persistence Query 本身并不是应用的查询侧数据库它的作用是帮助把数据从写侧迁移到查询侧数据库——你可以把它理解为写侧与读侧之间的数据搬运管道。在非常简单的场景下Persistence Query 的能力可能足以直接满足查询需求但在 CQRS 精神指引下官方强烈建议随着需求增长将写/读两侧拆分为各自独立优化的数据存储。对于 Durable State Behaviors持久化状态行为存在对等的查询接口实现见 Persistence Query using Durable State。想要把 Event Sourcing 与 CQRS 完整落地为可运行的应用程序可以学习官方的 Microservices with Akka 教程它演示了如何使用 Akka Persistence 与 Akka Projections 构建事件溯源 CQRS 应用。设计总览刻意保持宽松的 APIAkka Persistence Query故意被设计成一个非常宽松的 API 规范。这样做的目的是让通用 API 足够抽象使得每个 journal 实现都能暴露自己最强的能力例如 SQL journal 可以使用复杂 SQL 查询支持实时事件订阅的 journal 则可以暴露同样的 API——一个类型化的事件流。这带来一个关键约定每个 read journal 必须明确文档化它支持哪些查询类型。具体支持哪些查询、语义如何请查阅你所使用的 journal 插件文档。Akka Persistence Query 本身不提供任何 ReadJournal 的实际实现但它预定义了一批最常见的查询类型供大多数 journal 实现并非强制要求全部实现。真正的 ReadJournal 实现由社区插件提供每个插件针对特定数据存储如 Cassandra、JDBC 数据库。Read Journal查询的入口要发起查询首先需要通过PersistenceQuery扩展获取一个ReadJournal实例。ReadJournal 由社区插件实现每个插件针对特定数据存储并拥有一个插件标识符config path。例如对于提供akka.persistence.query.my-read-journal插件的库获取方式如下示例代码取自 PersistenceQueryDocSpec.scala 与 PersistenceQueryDocTest.java// Scala // obtain read journal by plugin id val readJournal PersistenceQuery(system).readJournalForMyScaladslReadJournal // issue query to journal val source: Source[EventEnvelope, NotUsed] readJournal.eventsByPersistenceId(user-1337, 0, Long.MaxValue) // materialize stream, consuming events source.runForeach { event println(Event: event) }// Java // obtain read journal by plugin id final MyJavadslReadJournal readJournal PersistenceQuery.get(system) .getReadJournalFor( MyJavadslReadJournal.class, akka.persistence.query.my-read-journal); // issue query to journal SourceEventEnvelope, NotUsed source readJournal.eventsByPersistenceId(user-1337, 0, Long.MAX_VALUE); // materialize stream, consuming events source.runForeach(event - System.out.println(Event: event), system);从源码看PersistenceQuery.scala 是一个标准的 Akka 扩展ExtensionId内部通过pluginFor(readJournalPluginId, ...)按插件标识符加载配置并委托给ReadJournalProvider分别产出 scaladsl 与 javadsl 两个版本的 ReadJournal。它同时支持传入一份自定义Config覆盖 ActorSystem 配置final def readJournalForT : scaladsl.ReadJournal: Tjournal 实现者被鼓励把插件标识符暴露为一个已知变量如NoopJournal.identifier方便用户通过readJournalForNoopJournal访问但这并非强制约定。预定义查询类型Akka Persistence Query 内置了若干查询接口并建议 journal 实现者按下述语义实现。需要注意这些查询类型虽然非常常见但 journal并没有义务全部实现——例如某种查询在特定 journal 中可能效率极低。使用前务必查阅你所用的 ReadJournal 插件文档确认其支持的具体查询类型与语义例如流何时完成。预定义查询分为以下几类每一类都提供live 流与current 快照流两个版本。PersistenceIdsQuery 与 CurrentPersistenceIdsQuerypersistenceIds用于订阅系统中全部持久化 ID的流。默认情况下该流是live流即 journal 应持续在系统出现新的 persistence id 时继续发射readJournal.persistenceIds() // Scalalive 流 readJournal.currentPersistenceIds() // Scala快照流如果使用场景不需要 live 流可以使用currentPersistenceIds——它只发射当前已存在的 ID到达末尾即完成。EventsByPersistenceIdQuery 与 CurrentEventsByPersistenceIdQueryeventsByPersistenceId在语义上等价于对某个 event sourced actor 做事件重放区别在于它返回的是流因此可以保持存活持续监视该persistenceId后续新持久化的事件readJournal.eventsByPersistenceId(user-us-1337, fromSequenceNr 0L, toSequenceNr Long.MaxValue)readJournal.eventsByPersistenceId(user-us-1337, 0L, Long.MAX_VALUE);大多数 journal 需要依赖轮询polling来实现这种live效果轮询间隔通常通过refresh-interval配置属性设置LevelDB 实现的默认值为3s见下文配置节。如果不需要 live 流使用currentEventsByPersistenceId即可。EventsByTag 与 CurrentEventsByTageventsByTag允许跨 persistenceId 查询事件——例如查询某个聚合根Aggregate Root类型的所有领域事件。这个查询在部分 journal 中实现困难或需要对数据存储做额外准备才能高效执行具体是否支持、如何支持请查阅 read journal 插件文档。事件的打标签由写侧完成可以使用事件溯源中的 tagging 机制也可以通过 Event Adapters 将事件包装为akka.persistence.journal.Tagged并携带tags。以 typedEventSourcedBehavior为例取自 BasicPersistentBehaviorCompileOnly.scalaval NumberOfEntityGroups 10 def tagEvent(entityId: String, event: Event): Set[String] { val entityGroup sgroup-${math.abs(entityId.hashCode % NumberOfEntityGroups)} event match { case _: OrderCompleted Set(entityGroup, order-completed) case _ Set(entityGroup) } } def apply(entityId: String): Behavior[Command] { EventSourcedBehaviorCommand, Event, State, emptyState State(), commandHandler (state, cmd) throw new NotImplementedError(TODO: process the command return an Effect), eventHandler (state, evt) throw new NotImplementedError(TODO: process the event return the next state)) .withTagger(event tagEvent(entityId, event)) }Java 对应版本见 BasicPersistentBehaviorTest.java。注意这里用withTagger把事件按实体组和业务类型如order-completed打上多个标签便于查询侧按标签订阅。查询侧用法Scala/Java// assuming journal is able to work with numeric offsets we can: val completedOrders: Source[EventEnvelope, NotUsed] readJournal.eventsByTag(order-completed, Offset.noOffset) // find first 10 completed orders: val firstCompleted: Future[Vector[OrderCompleted]] completedOrders .map(_.event) .collectType[OrderCompleted] .take(10) // cancels the query stream after pulling 10 elements .runFold(Vector.empty[OrderCompleted])(_ : _) // start another query, from the known offset val furtherOrders readJournal.eventsByTag(order-completed, offset Sequence(10))// assuming journal is able to work with numeric offsets we can: final SourceEventEnvelope, NotUsed completedOrders readJournal.eventsByTag(order-completed, new Sequence(0L)); // find first 10 completed orders: final CompletionStageListOrderCompleted firstCompleted completedOrders .map(EventEnvelope::event) .collectType(OrderCompleted.class) .take(10) // cancels the query stream after pulling 10 elements .runFold( new ArrayList(10), (acc, e) - { acc.add(e); return acc; }, system); // start another query, from the known offset SourceEventEnvelope, NotUsed furtherOrders readJournal.eventsByTag(order-completed, new Sequence(10));如示例所示查询流可以配合 Akka Streams 的所有常用算子比如take(10)取出前 10 个后取消流。内置EventsByTag查询还支持可选的Long类型 offset 参数journal 可以用它实现可续传流resumable stream例如用 SQL 的 WHERE 子句从指定行开始读取或者在能按插入时间排序事件的数据存储中把 Long 当作时间戳、只取更旧的事件。关于跨 persistenceId 查询有一个非常重要的注意事项事件在流中出现的顺序几乎不保证稳定即使在多次物化之间也不保证一致。journal可能选择提供严格排序但必须显式文档化其排序保证——例如按时间戳升序与 persistenceId 无关在关系型数据库上容易实现但在纯键值存储上则可能难以高效实现。NoOffset表示从该 tag 的第一个事件开始取Sequence(n)表示从第 n 个不含之后的事件开始取。每个事件在EventEnvelope中携带对应 offset可用于断点续传。不需要 live 流时使用currentEventsByTag。Offset 类型从 Offset.scala 源码可见offset 是抽象类目前定义了以下实现Offset 类型语义说明NoOffset从最开始读取取所有事件Java 侧用NoOffset.getInstance()Sequence(value: Long)有序序列号对应该 tag 的有序序号支持比较Ordered[Sequence]TimeBasedUUID(value: UUID)基于时间的 UUID构造时校验必须是 version 1 的 UUID否则抛IllegalArgumentExceptionTimestampOffset(timestamp, readTimestamp, seen)基于时间戳同一时间戳可能对应多个事件因此用seen: Map[String, Long]记录该时间戳下每个 persistenceId 已见过的序列号readTimestamp为微秒级读取时间TimestampOffsetBySlice(offsets: Map[Int, TimestampOffset])按 slice 的时间戳 offset每个 slice 一个TimestampOffset用于 typed 的按 slice 查询Offset对象还提供了便捷工厂Offset.noOffset、Offset.sequence(value)、Offset.timeBasedUUID(uuid)、Offset.timestamp(instant)。所有 offset 都是排他的——即流中不会包含与传入 offset 序号完全相同的事件因此可以直接把EventEnvelope中返回的 offset 作为下次查询的入参。EventEnvelope流中事件的统一包装所有查询流发射的元素都是 EventEnvelope它携带offset该事件在查询流中的位置用于续传persistenceId持久化该事件的 actor 标识sequenceNr该事件在对应persistenceId下的序列号event事件本体timestamp事件存储时间毫秒自 1970-01-01 UTC 起与System.currentTimeMillis一致metadata附加的元数据Scalametadata[M]/ JavagetMetadata(type)支持按类型提取与移除。其中persistenceId sequenceNr构成事件的唯一标识。EventsBySlice 与 CurrentEventsBySlicetyped APItyped 系列还提供了按slice查询的接口。slice 由 persistence id确定性计算得出目的是把所有 persistence id 均匀分布到各个 slice 上从而支持并行消费。相关接口见 typed/scaladsl/EventsBySliceQuery.scala 与 javadsl 对应文件trait EventsBySliceQuery extends ReadJournal { def eventsBySlicesEvent: Source[EventEnvelope[Event], NotUsed] def sliceForPersistenceId(persistenceId: String): Int def sliceRanges(numberOfRanges: Int): immutable.Seq[Range] }从源码注释看使用基于时间戳 offset 的EventsBySliceQuery实现还应当同时实现EventTimestampQuery与LoadEventQuery。其变体包括EventsBySliceStartingFromSnapshotsQuery/CurrentEventsBySliceStartingFromSnapshotsQuery从快照之后的事件开始查询减少重放量EventsBySliceFirehoseQuery当大量消费者例如同一实体类型的多个 Projection读取相同事件时它把来自数据库的事件流共享并扇出给各消费者流从而减少数据库查询与事件加载次数获得更好的扩展性。它通常与 Sharded Daemon Process 的 co-located 部署 配合使用。Firehose 的共享流行为可在 reference.conf 中配置关键项包括配置项默认值说明delegate-query-plugin-id必须由应用指定底层 EventsBySlice 查询插件的标识符broadcast-buffer-size256BroadcastHub 扇出缓冲大小必须是 2 的幂且小于 4096firehose-linger-timeout40s所有消费者关闭后共享流保留的时间之后再关闭catchup-overlap10s新消费者追赶共享流后的重叠窗口避免切换时丢事件deduplication-capacity10000重叠期间去重缓存的条目数slow-consumer-reaper-interval2s检测并中止慢消费者的后台任务间隔slow-consumer-lag-threshold5s判定慢消费者的延迟阈值abort-slow-consumer-after2s慢消费者持续该时长后被中止verbose-debug-loggingoff是否输出逐事件的调试日志生产慎开查询的物化值Materialized Valuesjournal 可以通过流框架的 物化值 特性在流物化时暴露额外信息。高级查询 journal 可以借此告知调用者所物化流的性质——例如流是有限还是无限、是否严格有序。物化值类型是返回Source的第二个类型参数这允许 journal 向用户提供专门的查询对象。示例如下取自测试代码// a plugin can provide: final case class RichEvent(tags: Set[String], payload: Any) case class QueryMetadata(deterministicOrder: Boolean, infinite: Boolean)static final class QueryMetadata { public final boolean deterministicOrder; public final boolean infinite; // ... }插件自定义查询byTagsWithMeta返回Source[RichEvent, QueryMetadata]使用方通过mapMaterializedValue读取元数据val query: Source[RichEvent, QueryMetadata] readJournal.byTagsWithMeta(Set(red, blue)) query .mapMaterializedValue { meta println( sThe query is: sordered deterministically: ${meta.deterministicOrder}, sinfinite: ${meta.infinite}) } .map { event println(sEvent payload: ${event.payload}) } .runWith(Sink.ignore)性能与反规范化从写侧投影到读侧使用 Event Sourcing 与 CQRS 构建系统时必须意识到写侧与读侧的诉求完全不同把两者分离到各自优化的数据存储中才能让两侧都获得最佳体验。举例在竞价bidding系统中写侧要尽快落盘并回复出价人因此写吞吐量优先级最高——这往往意味着具备高扩展写能力的数据存储其查询表达能力反而较弱。而同一应用可能还需要复杂的统计视图或分析师要基于数据寻找最优竞价策略——这通常需要 SQL 这类表达力强的查询甚至写 Spark 作业来分析数据。因此写侧存储的数据必须被**投影projected**到另一个读优化的数据存储中。Akka Persistence 语境下的物化视图Materialized View指某次查询结果的持久化存储。也就是说视图只创建一次之后被反复查询——以这种形式查询比直接查询源事件更高效或更有意义。方式一投影到 Reactive Streams 兼容的数据存储如果读侧数据存储暴露了 Reactive Streams 接口实现一个简单投影就是把 read journal 的流直接喂给数据库驱动接口示例取自 PersistenceQueryDocSpec.scala 与 PersistenceQueryDocTest.javaval readJournal PersistenceQuery(system).readJournalForMyScaladslReadJournal val dbBatchWriter: Subscriber[immutable.Seq[Any]] ReactiveStreamsCompatibleDBDriver.batchWriter // Using an example (Reactive Streams) Database driver readJournal .eventsByPersistenceId(user-1337, fromSequenceNr 0L, toSequenceNr Long.MaxValue) .map(envelope envelope.event) .map(convertToReadSideTypes) // convert to datatype .grouped(20) // batch inserts into groups of 20 .runWith(Sink.fromSubscriber(dbBatchWriter)) // write batches to read-side databasefinal ReactiveStreamsCompatibleDBDriver driver new ReactiveStreamsCompatibleDBDriver(); final SubscriberListObject dbBatchWriter driver.batchWriter(); // Using an example (Reactive Streams) Database driver readJournal .eventsByPersistenceId(user-1337, 0L, Long.MAX_VALUE) .map(envelope - envelope.event()) .grouped(20) // batch inserts into groups of 20 .runWith(Sink.fromSubscriber(dbBatchWriter), system); // write batches to read-side database方式二使用 mapAsync 投影如果目标数据库没有提供可执行写入的 Reactive StreamsSubscriber则只能用普通函数或 Actor 实现写入逻辑。当写逻辑无状态、只需把事件从一种数据类型转换为另一种再写入时投影长这样trait ExampleStore { def save(event: Any): Future[Unit] } val store: ExampleStore ??? readJournal .eventsByTag(bid, NoOffset) .mapAsync(1) { e store.save(e) } .runWith(Sink.ignore)static class ExampleStore { CompletionStageVoid save(Object any) { /* ... */ } } final ExampleStore store new ExampleStore(); readJournal .eventsByTag(bid, new Sequence(0L)) .mapAsync(1, store::save) .runWith(Sink.ignore(), system);可续传投影Resumable Projections某些场景需要可续传的投影每次运行不必从起点重新处理而是把已处理事件的序列号offset保存下来下次启动时从该 offset 继续。这个模式已经由Akka Projections模块实现推荐直接使用避免重复造轮子。查询插件Query Plugins开发指南查询插件是面向各种数据存储的ReadJournal实现绝大多数由社区维护完整列表见 Akka Persistence Query 的 Community Plugins 页面。本节为需要自定义插件的开发者提供指引——大多数用户无需自己实现 journal除非目标数据存储尚无支持。由于不同数据存储提供的查询能力差异巨大journal 插件必须详尽文档化其暴露的语义及处理的查询场景。ReadJournalProvider API一个 read journal 插件必须实现 ReadJournalProvider它分别创建 scaladsl 与 javadsl 两个版本的ReadJournal。插件必须同时实现两个 DSL因为akka.stream.scaladsl.Source与akka.stream.javadsl.Source是不同类型——虽然两者可以互相转换但让最终用户直接拿到对应语言的Source最为方便。如下所示其中一个实现可以委托给另一个。一个极简的插件实现完整版见 PersistenceQueryDocSpec.scala 与 PersistenceQueryDocTest.javaclass MyReadJournalProvider(system: ExtendedActorSystem, config: Config) extends ReadJournalProvider { private val readJournal: MyScaladslReadJournal new MyScaladslReadJournal(system, config) override def scaladslReadJournal(): MyScaladslReadJournal readJournal override def javadslReadJournal(): MyJavadslReadJournal new MyJavadslReadJournal(readJournal) } class MyScaladslReadJournal(system: ExtendedActorSystem, config: Config) extends akka.persistence.query.scaladsl.ReadJournal with akka.persistence.query.scaladsl.EventsByTagQuery with akka.persistence.query.scaladsl.EventsByPersistenceIdQuery with akka.persistence.query.scaladsl.PersistenceIdsQuery with akka.persistence.query.scaladsl.CurrentPersistenceIdsQuery { private val refreshInterval: FiniteDuration config.getDuration(refresh-interval, MILLISECONDS).millis override def eventsByTag(tag: String, offset: Offset): Source[EventEnvelope, NotUsed] offset match { case Sequence(offsetValue) Source.fromGraph(new MyEventsByTagSource(tag, offsetValue, refreshInterval)) case NoOffset eventsByTag(tag, Sequence(0L)) //recursive case _ throw new IllegalArgumentException(MyJournal does not support offset.getClass.getName offsets) } override def eventsByPersistenceId( persistenceId: String, fromSequenceNr: Long, toSequenceNr: Long): Source[EventEnvelope, NotUsed] { // implement in a similar way as eventsByTag ??? } override def persistenceIds(): Source[String, NotUsed] ??? override def currentPersistenceIds(): Source[String, NotUsed] ??? }其中eventsByTag可以基于一个自定义GraphStage实现仓库测试中提供了参考实现 MyEventsByTagSource.scalaJava 版见 MyEventsByTagSource.java。ReadJournalProvider类必须提供以下签名之一的构造器ReadJournalProvider.scala带ExtendedActorSystem、com.typesafe.config.Config、String配置路径三个参数带ExtendedActorSystem与com.typesafe.config.Config两个参数只带一个ExtendedActorSystem参数无参构造器。ActorSystem 配置中该插件 section 的配置会被传入 config 构造参数插件的配置路径通过String参数传入。测试代码中的插件配置形如akka.persistence.query.my-read-journal { class docs.persistence.query.PersistenceQueryDocSpec$MyReadJournalProvider refresh-interval 3s }如果底层数据存储只支持到达结果集末尾即完成的查询那么为了支持无限事件流包含初始查询完成之后新存储的事件journal 必须每隔一段时间重新提交查询。建议插件使用名为refresh-interval的配置属性来定义这个刷新间隔。参考实现LevelDB 查询插件仓库内置了 LevelDB journal 的查询实现其文档见 persistence-query-leveldb.md代码位于 akka-persistence-query/src/main/scala/akka/persistence/query/journal/leveldb含AllPersistenceIdsStage、EventsByPersistenceIdStage、EventsByTagStage等流阶段。LevelDB 实现目前已弃用官方建议新应用改用 Akka Persistence JDBC。其默认配置reference.confakka.persistence.query.journal.leveldb { # Implementation class of the LevelDB ReadJournalProvider class akka.persistence.query.journal.leveldb.LeveldbReadJournalProvider # Absolute path to the write journal plugin configuration entry that this # query journal will connect to. That must be a LeveldbJournal or SharedLeveldbJournal. # If undefined (or ) it will connect to the default journal as specified by the # akka.persistence.journal.plugin property. write-plugin # The LevelDB write journal is notifying the query side as soon as things # are persisted, but for efficiency reasons the query side retrieves the events # in batches that sometimes can be delayed up to the configured refresh-interval. refresh-interval 3s # How many events to fetch in one query (replay) and keep buffered until they # are delivered downstreams. max-buffer-size 100 }配置项默认值说明classLeveldbReadJournalProvider插件实现类write-plugin关联的写 journal 插件配置绝对路径留空则使用akka.persistence.journal.plugin指定的默认 journalrefresh-interval3s查询侧批量拉取事件的延迟上限max-buffer-size100单次查询重放拉取并缓冲的事件数LevelDB 实现的语义值得借鉴也可作为理解各插件文档的模板eventsByPersistenceId按序列号排序与写入顺序一致多次执行的流前缀同顺序相同除非事件被删除流不因到达当前事件末尾而完成会持续推送新事件currentEventsByPersistenceId才会完成查询后端失败时流以失败结束。persistenceIds流无序多次执行顺序可能不同新 actor 创建后持续推送新 ID且该查询没有轮询与批处理写 journal 立即通知查询侧currentPersistenceIds到达末尾即完成。eventsByTag按 offsettag 序列号排序NoOffset取全部Sequence从指定位置续传offset 排他使用deleteMessages(toSequenceNr)删除的事件不会从 tagged 流中消失persistenceId sequenceNr是事件唯一标识。横向扩展Scaling Out当事件数量极大、每个事件的处理工作量很高或者对韧性resilience有要求——节点崩溃后持久化查询能快速在新节点上启动并恢复——时结合事件打标签event tagging与 Cluster Sharding是跨集群分片事件的绝佳方案事件按实体打上标签后通过 Cluster Sharding 把不同分片的事件分散到集群各节点并行处理任一节点故障时其分片可被重分配到其他节点继续消费。示例项目官方的 Microservices with Akka 教程演示了如何将 Event Sourcing 与 Projections 结合使用事件被标记tagged供事件处理器消费以构建事件的其他表示形式或将事件发布给其他服务。这是 Persistence Query 在真实微服务架构中的典型落地路径——写侧持久化事件并打标签读侧通过查询流构建投影与派生数据。总结Akka Persistence Query 为 Akka 生态提供了统一、异步、基于流的事件查询抽象。使用时要牢记三个关键点每个 journal 插件自行声明支持的查询类型与语义务必查阅插件文档跨 persistenceId 的查询如eventsByTag不保证稳定顺序续传依赖EventEnvelope中携带的排他性 offset真正的 CQRS 读侧建议将事件投影到独立优化的数据存储物化视图简单场景可直接使用查询流复杂场景则配合 Akka Projections 与 Cluster Sharding 实现可续传、可扩展的投影。【免费下载链接】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),仅供参考