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-core
Akka 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> <groupId>com.typesafe.akka</groupId> <artifactId>akka-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 插件实现,以暴露各自的查询能力。
其最典型的应用场景是 CQRS(Command 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 应用。
设计总览:刻意保持"宽松"的 API
Akka 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 Source<EventEnvelope, 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 与 CurrentPersistenceIdsQuery
persistenceIds用于订阅系统中全部持久化 ID的流。默认情况下该流是"live"流,即 journal 应持续在系统出现新的 persistence id 时继续发射:
readJournal.persistenceIds() // Scala,live 流 readJournal.currentPersistenceIds() // Scala,快照流如果使用场景不需要 live 流,可以使用currentPersistenceIds——它只发射当前已存在的 ID,到达末尾即完成。
EventsByPersistenceIdQuery 与 CurrentEventsByPersistenceIdQuery
eventsByPersistenceId在语义上等价于对某个 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 与 CurrentEventsByTag
eventsByTag允许跨 persistenceId 查询事件——例如查询某个聚合根(Aggregate Root)类型的所有领域事件。这个查询在部分 journal 中实现困难,或需要对数据存储做额外准备才能高效执行;具体是否支持、如何支持,请查阅 read journal 插件文档。
事件的"打标签"由写侧完成:可以使用事件溯源中的 tagging 机制,也可以通过 Event Adapters 将事件包装为akka.persistence.journal.Tagged并携带tags。以 typedEventSourcedBehavior为例(取自 BasicPersistentBehaviorCompileOnly.scala):
val NumberOfEntityGroups = 10 def tagEvent(entityId: String, event: Event): Set[String] = { val entityGroup = s"group-${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 Source<EventEnvelope, NotUsed> completedOrders = readJournal.eventsByTag("order-completed", new Sequence(0L)); // find first 10 completed orders: final CompletionStage<List<OrderCompleted>> 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 Source<EventEnvelope, 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,否则抛IllegalArgumentException |
TimestampOffset(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 与 CurrentEventsBySlice(typed API)
typed 系列还提供了按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-size | 256 | BroadcastHub 扇出缓冲大小,必须是 2 的幂且小于 4096 |
firehose-linger-timeout | 40s | 所有消费者关闭后共享流保留的时间,之后再关闭 |
catchup-overlap | 10s | 新消费者追赶共享流后的重叠窗口,避免切换时丢事件 |
deduplication-capacity | 10000 | 重叠期间去重缓存的条目数 |
slow-consumer-reaper-interval | 2s | 检测并中止慢消费者的后台任务间隔 |
slow-consumer-lag-threshold | 5s | 判定慢消费者的延迟阈值 |
abort-slow-consumer-after | 2s | 慢消费者持续该时长后被中止 |
verbose-debug-logging | off | 是否输出逐事件的调试日志(生产慎开) |
查询的物化值(Materialized Values)
journal 可以通过流框架的 物化值 特性,在流物化时暴露额外信息。高级查询 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( s"The query is: " + s"ordered deterministically: ${meta.deterministicOrder}, " + s"infinite: ${meta.infinite}") } .map { event => println(s"Event 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.java):
val 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 Subscriber<List<Object>> 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 { CompletionStage<Void> 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.java):
class 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.scala(Java 版见 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.conf):
akka.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 }| 配置项 | 默认值 | 说明 |
|---|---|---|
class | LeveldbReadJournalProvider | 插件实现类 |
write-plugin | "" | 关联的写 journal 插件配置绝对路径;留空则使用akka.persistence.journal.plugin指定的默认 journal |
refresh-interval | 3s | 查询侧批量拉取事件的延迟上限 |
max-buffer-size | 100 | 单次查询(重放)拉取并缓冲的事件数 |
LevelDB 实现的语义值得借鉴,也可作为理解各插件文档的模板:
eventsByPersistenceId:按序列号排序,与写入顺序一致;多次执行的流前缀(同顺序)相同(除非事件被删除);流不因到达当前事件末尾而完成,会持续推送新事件(currentEventsByPersistenceId才会完成);查询后端失败时流以失败结束。persistenceIds:流无序,多次执行顺序可能不同;新 actor 创建后持续推送新 ID,且该查询没有轮询与批处理(写 journal 立即通知查询侧);currentPersistenceIds到达末尾即完成。eventsByTag:按 offset(tag 序列号)排序;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),仅供参考