- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
本指南以 Akka 官方文档 persistence-query-leveldb.md 为骨架,围绕akka-persistence-query模块中 LevelDB 读日志(ReadJournal)插件的使用展开:从依赖引入、ReadJournal获取、六种查询 API 的语义与用法,到reference.conf配置项与底层 Stage 实现原理,再到测试与弃用注意事项。读完本文,你将能够在 Akka 应用中基于 LevelDB 写出可运行的事件溯源查询代码,并理解其"写日志推送 + 批量刷新"的底层工作方式,为迁移到其他日志实现(如 JDBC)打好基础。
1. 现状与定位:LevelDB 日志查询插件已弃用
官方文档开篇即给出明确结论:LevelDB 日志与查询插件已弃用(deprecated),不建议在新应用中使用,官方推荐的替代方案是 Akka Persistence JDBC。在源码中也能看到这一标记——scaladsl/LeveldbReadJournal.scala 与 javadsl/LeveldbReadJournal.scala 的类声明上都标注了:
@deprecated("Use another journal implementation", "2.6.15")即自 2.6.15 起标记弃用。尽管如此,该插件依然是理解 Akka Persistence Query 抽象模型(ReadJournal、Offset、EventEnvelope、Tagged 等)最直观的入门样例,且测试代码与文档仍然保留,可作为学习资料。
注意:本文所有内容均基于当前仓库(akka-core)中保留的 LevelDB 实现,若用于生产环境,请优先评估官方推荐的 JDBC 等替代方案。
2. 添加依赖
使用 Persistence Query 需要在项目中引入akka-persistence-query模块,它会同时传递依赖akka-persistence模块(详见 persistence.md)。
sbt
libraryDependencies += "com.typesafe.akka" %% "akka-persistence-query" % AkkaVersionMaven
<properties> <akka.version>2.9.x</akka.version> <scala.binary.version>2.13</scala.binary.version> </properties> <dependency> <groupId>com.typesafe.akka</groupId> <artifactId>akka-persistence-query_${scala.binary.version}</artifactId> <version>${akka.version}</version> </dependency>Gradle
def versions = [ ScalaBinary: "2.13" ] def akkaVersion = "2.9.x" dependencies { implementation "com.typesafe.akka:akka-persistence-query_${versions.ScalaBinary}:${akkaVersion}" }其中AkkaVersion请替换为实际使用的 Akka 版本号。文档同时提醒:Akka 依赖可通过 Akka 的 secure library repository 获取,访问时可能需要使用带 token 的安全 URL(详见 Akka 官方账号页面说明)。另外,LevelDB 写日志插件本身还要求显式引入 LevelDB 实现依赖(见 akka-persistence 的 reference.conf 中的注释),可选用org.iq80.leveldb(Java 移植版)或org.fusesource.leveldbjni(JNI 原生版)。
3. 如何获取 ReadJournal
ReadJournal通过akka.persistence.query.PersistenceQuery扩展(extension)获取。核心要点:
- 使用
LeveldbReadJournal.Identifier(常量值为"akka.persistence.query.journal.leveldb")作为插件标识; - 该标识同时是配置文件中的绝对路径,二者一一对应。
Scala(完整代码见 LeveldbPersistenceQueryDocSpec.scala):
import akka.persistence.query.PersistenceQuery import akka.persistence.query.journal.leveldb.scaladsl.LeveldbReadJournal val queries = PersistenceQuery(system).readJournalForLeveldbReadJournalJava(完整代码见 LeveldbPersistenceQueryDocTest.java):
import akka.persistence.query.PersistenceQuery; import akka.persistence.query.journal.leveldb.javadsl.LeveldbReadJournal; LeveldbReadJournal queries = PersistenceQuery.get(system) .getReadJournalFor(LeveldbReadJournal.class, LeveldbReadJournal.Identifier());从源码看,PersistenceQuery(system).readJournalFor[...]会根据配置中的class字段实例化 LeveldbReadJournalProvider(其class = "akka.persistence.query.journal.leveldb.LeveldbReadJournalProvider"),再由 Provider 分别创建 Scala DSL 与 Java DSL 两个门面对象。所有查询方法最终返回akka.stream.scaladsl.Source类型的 Akka Streams 数据流。
注意:本文档描述的仅是LevelDB 这一种日志实现的 Persistence Query 语义。官方 persistence-query.md 也明确说明:不同日志实现的查询语义可能不同,例如其他实现可能不支持全部查询类型、offset 类型或排序保证,使用时务必查阅对应实现的文档。
4. 支持的查询:六种 API 全面解析
LeveldbReadJournal实现了 Persistence Query 中的六种查询接口(源码见 scaladsl/LeveldbReadJournal.scala 的类声明):
class LeveldbReadJournal(system: ExtendedActorSystem, config: Config) extends ReadJournal with PersistenceIdsQuery with CurrentPersistenceIdsQuery with EventsByPersistenceIdQuery with CurrentEventsByPersistenceIdQuery with EventsByTagQuery with CurrentEventsByTagQuery下面按三组逐一讲解。
4.1 EventsByPersistenceId 与 CurrentEventsByPersistenceId
eventsByPersistenceId用于检索某个指定persistenceId的PersistentActor所持久化的全部事件。
Scala:
val queries = PersistenceQuery(system).readJournalForLeveldbReadJournal val src: Source[EventEnvelope, NotUsed] = queries.eventsByPersistenceId("some-persistence-id", 0L, Long.MaxValue) val events: Source[Any, NotUsed] = src.map(_.event)Java:
LeveldbReadJournal queries = PersistenceQuery.get(system) .getReadJournalFor(LeveldbReadJournal.class, LeveldbReadJournal.Identifier()); Source<EventEnvelope, NotUsed> source = queries.eventsByPersistenceId("some-persistence-id", 0, Long.MAX_VALUE);语义要点(文档明确约定,源码注释亦有完整复述):
- 区间检索:可通过
fromSequenceNr与toSequenceNr限定事件子集;使用0L和Long.MaxValue(Java 为Long.MAX_VALUE)则可检索全部事件。每个事件的序号会体现在EventEnvelope中,因此可以"从某个序号之后继续",实现断点续流。 - 顺序保证:返回流按序号(sequence number)排序,即与
PersistentActor持久化事件的顺序一致。同一查询多次执行返回相同前缀、相同顺序的流元素,除非事件被删除。 - Live 与 Current 的区别:
eventsByPersistenceId是"活流"——到达当前已存事件末尾时不会结束,而是持续推送新持久化的事件;currentEventsByPersistenceId则是"当前快照流"——到达已存事件末尾即完成。 - 批处理刷新:LevelDB 写日志会在事件持久化后尽快通知查询端,但出于效率考虑,查询端按批次拉取事件,最多可能延迟到配置的
refresh-interval(或显式传入的RefreshInterval提示)时长。 - 失败语义:若后端日志执行查询失败,流将以 failure 结束。
4.2 PersistenceIds 与 CurrentPersistenceIds
persistenceIds用于检索所有持久化 Actor 的persistenceId列表。
Scala:
val queries = PersistenceQuery(system).readJournalForLeveldbReadJournal val src: Source[String, NotUsed] = queries.persistenceIds()Java:
LeveldbReadJournal queries = PersistenceQuery.get(system) .getReadJournalFor(LeveldbReadJournal.class, LeveldbReadJournal.Identifier()); Source<String, NotUsed> source = queries.persistenceIds();语义要点:
- 无序:返回流不保证顺序,多次执行可能得到不同的排列顺序。
- Live 与 Current 的区别:
persistenceIds是活流,新持久化 Actor 出现时会持续推送新的persistenceId;currentPersistenceIds则在遍历完当前集合后完成。 - 无轮询:与其他两种查询不同,
persistenceIds不涉及周期轮询或批量拉取——写日志一旦创建新的persistenceId会立刻通知查询端。这在源码中有直接体现(scaladsl/LeveldbReadJournal.scala 注释"no polling for this query, the write journal will push all changes, i.e. no refreshInterval"),其底层由 AllPersistenceIdsStage 实现。 - 失败语义:后端日志执行失败时流以 failure 结束。
4.3 EventsByTag 与 CurrentEventsByTag
eventsByTag用于检索带有指定 tag 的事件,典型场景是"某个聚合根(Aggregate Root)类型的所有领域事件"。
Scala:
val queries = PersistenceQuery(system).readJournalForLeveldbReadJournal val src: Source[EventEnvelope, NotUsed] = queries.eventsByTag(tag = "green", offset = Sequence(0L))Java:
LeveldbReadJournal queries = PersistenceQuery.get(system) .getReadJournalFor(LeveldbReadJournal.class, LeveldbReadJournal.Identifier()); Source<EventEnvelope, NotUsed> source = queries.eventsByTag("green", new Sequence(0L));如何给事件打标签:WriteEventAdapter + Tagged
要给事件打上 tag,需要创建一个 Event Adapters(写事件适配器),把事件包装进akka.persistence.journal.Tagged,并附上tags集合。
Scala(MyTaggingEventAdapter,摘自 LeveldbPersistenceQueryDocSpec.scala):
import akka.persistence.journal.WriteEventAdapter import akka.persistence.journal.Tagged class MyTaggingEventAdapter extends WriteEventAdapter { val colors = Set("green", "black", "blue") override def toJournal(event: Any): Any = event match { case s: String => val tags = colors.foldLeft(Set.empty[String]) { (acc, c) => if (s.contains(c)) acc + c else acc } if (tags.isEmpty) event else Tagged(event, tags) case _ => event } override def manifest(event: Any): String = "" }Java(摘自 LeveldbPersistenceQueryDocTest.java):
import akka.persistence.journal.Tagged; import akka.persistence.journal.WriteEventAdapter; import java.util.HashSet; import java.util.Set; public static class MyTaggingEventAdapter implements WriteEventAdapter { @Override public Object toJournal(Object event) { if (event instanceof String) { String s = (String) event; Set<String> tags = new HashSet<String>(); if (s.contains("green")) tags.add("green"); if (s.contains("black")) tags.add("black"); if (s.contains("blue")) tags.add("blue"); if (tags.isEmpty()) return event; else return new Tagged(event, tags); } else { return event; } } @Override public String manifest(Object event) { return ""; } }该适配器把包含"green"、"black"、"blue"等颜色关键词的字符串事件分别打上对应 tag;无匹配时不包装、原样返回。之后还需要在配置中把该适配器绑定到具体事件类型(绑定方式见下文第 6 节的测试配置示例event-adapters/event-adapter-bindings)。
Offset 语义(重要)
- Offset 类型:
eventsByTag支持NoOffset(检索该 tag 的全部事件)或Sequence类型 offset(检索子集)。Sequenceoffset 对应"该 tag 维度的有序序号"。源码中(scaladsl/LeveldbReadJournal.scala)对 offset 做了严格匹配:Sequence正常处理、NoOffset递归等价于Sequence(0L),其他 offset 类型直接抛出IllegalArgumentException("LevelDB does not support ... offsets")——即 LevelDB 实现不支持TimeBasedUUID等其他 offset。 - Offset 是排他的:与 offset 序号完全相同的那条事件不会被包含在返回流中。这意味着你可以把
EventEnvelope中返回的 offset 直接作为下一次查询的offset参数,实现精确续传、不重不漏。 - Envelope 附加信息:除 offset 外,
EventEnvelope还提供persistenceId与sequenceNr。其中sequenceNr是持久化该事件的 Actor 自己的序号,persistenceId+sequenceNr构成事件的唯一标识。 - 排序与稳定性:返回流按 offset(tag 序号)排序,与写日志存储顺序一致;多次执行返回相同元素、相同顺序。
- 删除不影响 tag 流:文档用专门的 note 强调——通过
deleteMessages(toSequenceNr)删除的事件不会从 "tagged stream"(tag 事件流)中删除。
Live 与 Current 的区别
与前面两组一致:eventsByTag是活流,到达当前已存事件末尾后继续推送新事件;currentEventsByTag到达末尾即完成。eventsByTag同样采用"写日志推送 + 按refresh-interval批量拉取"的模式,后端失败时流以 failure 结束。
5. 三种查询的实现机制对比
| 查询 | 排序 | 多次执行稳定性 | 到达末尾行为 | 通知/拉取模式 |
|---|---|---|---|---|
eventsByPersistenceId | 按序号升序 | 稳定(除非事件被删除) | live:继续推送;current:完成 | 写日志推送 + 按refresh-interval批量拉取 |
persistenceIds | 无序 | 不保证 | live:继续推送;current:完成 | 写日志即时推送,无轮询无批量 |
eventsByTag | 按 tag offset 升序 | 稳定 | live:继续推送;current:完成 | 写日志推送 + 按refresh-interval批量拉取 |
从源码结构看(可推断):三个查询分别由 EventsByPersistenceIdStage、AllPersistenceIdsStage、EventsByTagStage 三个自定义 GraphStage 驱动,配合 Buffer 实现背压与批量缓存;current*变体与 live 变体共用同一 Stage,只是不传入refreshInterval(传None)并把liveQuery置为false。
6. 配置详解
LevelDB 读日志的配置项位于绝对路径"akka.persistence.query.journal.leveldb"下(与LeveldbReadJournal.Identifier一致)。完整默认配置见 akka-persistence-query 的 reference.conf:
# Configuration for the LeveldbReadJournal 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 | akka.persistence.query.journal.leveldb.LeveldbReadJournalProvider | ReadJournal 的 Provider 实现类,一般无需修改 |
write-plugin | ""(空字符串) | 所连接的写日志插件配置的绝对路径,必须是LeveldbJournal或SharedLeveldbJournal;留空则使用akka.persistence.journal.plugin指定的默认日志 |
refresh-interval | 3s | 写日志在事件持久化后立即通知查询端,但查询端按批次拉取,最多可延迟到该时长(作用于eventsByPersistenceId与eventsByTag) |
max-buffer-size | 100 | 单次查询(replay)拉取并缓存在下游投递之前的事件数量 |
6.1 write-plugin 的解析与校验(源码级)
从 scaladsl/LeveldbReadJournal.scala 可以看到这三个配置项在构造时的实际消费逻辑:
private val refreshInterval = Some(config.getDuration("refresh-interval", MILLISECONDS).millis) private val writeJournalPluginId: String = config.getString("write-plugin") private val maxBufSize: Int = config.getInt("max-buffer-size") private val resolvedWriteJournalPluginId = if (writeJournalPluginId.isEmpty) system.settings.config.getString("akka.persistence.journal.plugin") else writeJournalPluginId require( resolvedWriteJournalPluginId.nonEmpty && system.settings.config .getConfig(resolvedWriteJournalPluginId) .getString("class") == "akka.persistence.journal.leveldb.LeveldbJournal", s"Leveldb read journal can only work with a Leveldb write journal. Current plugin [$resolvedWriteJournalPluginId] is not a LeveldbJournal")要点:
write-plugin为空时,自动回退到akka.persistence.journal.plugin指定的默认日志;- 构造时会
require校验:解析出的写日志插件配置中class必须恰为akka.persistence.journal.leveldb.LeveldbJournal,否则直接抛出IllegalArgumentException("Leveldb read journal can only work with a Leveldb write journal")。这是读日志与写日志必须配套的强约束——即使write-plugin配的是SharedLeveldbJournal,读端校验依旧按此执行。 refreshInterval只在 live 查询中传入 Stage;current*查询传None(见上文实现机制对比)。
6.2 配套的 LevelDB 写日志配置
读日志需要与写日志配套使用。写日志的默认配置位于 akka-persistence 的 reference.conf:
# LevelDB journal plugin. # Note: this plugin requires explicit LevelDB dependency, see below. akka.persistence.journal.leveldb { class = "akka.persistence.journal.leveldb.LeveldbJournal" plugin-dispatcher = "akka.persistence.dispatchers.default-plugin-dispatcher" replay-dispatcher = "akka.persistence.dispatchers.default-replay-dispatcher" dir = "journal" # Storage location of LevelDB files. fsync = on # Use fsync on write. checksum = off # Verify checksum on read. native = on # Native LevelDB (via JNI) or LevelDB Java port. compaction-intervals { # Number of deleted messages per persistence id that will trigger compaction } }另有仅供测试使用的akka.persistence.journal.leveldb-shared(SharedLeveldbJournal)。测试环境通常这样配置(摘自 EventsByTagSpec.scala 的测试配置):
akka.persistence.journal.plugin = "akka.persistence.journal.leveldb" akka.persistence.journal.leveldb { dir = "target/journal-EventsByTagSpec" event-adapters { color-tagger = akka.persistence.query.journal.leveldb.ColorTagger } event-adapter-bindings { "java.lang.String" = color-tagger } } akka.persistence.query.journal.leveldb { refresh-interval = 1s max-buffer-size = 2 }其中event-adapters/event-adapter-bindings正是把上一节编写的WriteEventAdapter(如MyTaggingEventAdapter、ColorTagger)绑定到具体事件类型(如java.lang.String)的配置方式——只写适配器类而不做绑定,tag 不会生效。测试还演示了如何通过 HOCON 变量替换派生一个几乎相同的配置副本(leveldb-no-refresh = ${akka.persistence.query.journal.leveldb}并覆盖refresh-interval = 10m),用于验证长刷新间隔下的语义。
7. 底层实现原理:写日志如何"通知"查询端
Persistence Query for LevelDB 的核心设计是事件推送式:查询端并不持续轮询 LevelDB 文件,而是依赖写日志(LeveldbJournal,源码见 akka-persistence/src/main/scala/akka/persistence/journal/leveldb/LeveldbJournal.scala)在事件持久化后主动向查询端发出通知。结合源码结构可以归纳出三条链路:
- 事件流(eventsByPersistenceId / eventsByTag):写日志持久化新事件后推送通知 → 查询端收到通知后按
max-buffer-size批量执行 replay 拉取 → 数据在Buffer中排队,按下游背压逐个投递;批量拉取可能最多延迟refresh-interval。 - ID 流(persistenceIds):写日志在每次创建新的
persistenceId时即时推送,无轮询、无批处理,因此延迟最低(源码注释明确"there is no periodic polling or batching involved in this query")。 - tag 流(eventsByTag):与事件流机制相同,但检索维度是 tag 序号(offset),且删除(
deleteMessages)不会影响该流——tag 事件一经写入即长期保留在流中。
这三个查询各自对应的EventsByPersistenceIdStage、AllPersistenceIdsStage、EventsByTagStage都位于 akka-persistence-query/src/main/scala/akka/persistence/query/journal/leveldb/ 目录下,与Buffer.scala共同构成 LevelDB 读日志的运行时。仓库中对应的测试套件(EventsByPersistenceIdSpec.scala、AllPersistenceIdsSpec.scala、EventsByTagSpec.scala)覆盖了活流推送、current 完成、offset 排他续传、删除不影响 tag 流等全部文档语义,是理解本文各语义要点最直接的验证入口。
8. 实战:从查询到投影的最小闭环
把上述要素组合起来,一个基于 LevelDB 的典型查询/投影场景包含四步:
- 引入依赖:
akka-persistence-query+ LevelDB 实现(见第 2 节); - 配置写日志与读日志:设置
akka.persistence.journal.plugin指向 LevelDB 写日志,按需覆盖akka.persistence.query.journal.leveldb的refresh-interval、max-buffer-size; - 编写并绑定 Tag 适配器(如需
eventsByTag):实现WriteEventAdapter返回Tagged,并在event-adapters/event-adapter-bindings中绑定; - 获取 ReadJournal 并消费流:通过
PersistenceQuery(system).readJournalForLeveldbReadJournal拿到查询对象,对返回的Source施加map、filter、runWith等任意 Akka Streams 操作。
例如,从Sequence(0L)开始消费 tag 为"green"的事件,并把每个EventEnvelope中的 offset 保存下来,下次查询时直接传入该 offset(因为 offset 排他,事件不会重复消费):
val queries = PersistenceQuery(system).readJournalForLeveldbReadJournal queries .eventsByTag(tag = "green", offset = savedOffset) // savedOffset 来自上一次查询的 envelope.offset .map { env => (env.persistenceId, env.sequenceNr, env.offset, env.event) } .runForeach { case (pid, seqNr, offset, event) => // 处理事件,并记录 offset 以便断点续传 }9. 注意事项与迁移建议
- 弃用声明:LevelDB 读日志自 Akka 2.6.15 起标记
@deprecated,新项目请勿选用;存量项目建议评估迁移到 Akka Persistence JDBC 等受支持实现。不同实现的查询语义(offset 类型、排序保证、删除语义)可能不同,迁移时必须对照目标实现的文档逐项核对。 - offset 排他性:
eventsByTag的Sequenceoffset 是排他的,务必使用EventEnvelope.offset作为续传游标,不要自行offset + 1(不同实现可能不同)。 - LevelDB 仅支持 Sequence/NoOffset:传入其他 offset 类型(如
TimeBasedUUID)会直接抛IllegalArgumentException,这是源码层面的硬约束。 - 写读必须配套:读日志构造时会校验写日志插件
class必须是LeveldbJournal,否则启动即失败。 - 延迟特性:
eventsByPersistenceId/eventsByTag存在最多一个refresh-interval的批量拉取延迟;persistenceIds则无延迟、无轮询。若对时效敏感,可调小refresh-interval,但会以更多次批量 replay 为代价。 - 删除语义差异:
eventsByPersistenceId中已删除事件不再出现,而eventsByTag中已删除事件依然保留在流中,设计消费逻辑时需注意两者的不对称性。
延伸阅读:Persistence Query 通用 API 与各日志实现差异见 persistence-query.md;事件适配器与Tagged的完整机制见 persistence.md;LevelDB 写日志实现见 LeveldbJournal.scala 及其配套的 LeveldbStore.scala。
- 后端
- 并发编程
- 异步编程
【免费下载链接】akka-core
A platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.
相关推荐
Akka Persistence Query 完全指南:基于 Durable State 的 CQRS 查询端实战
Akka Persistence Query 完全指南:基于 Durable State 的 CQRS 查询端实战 导读 本文聚焦 Akka 中面向 Durab
后端并发编程异步编程Akka Persistence Query 实战指南:用统一异步流接口构建 CQRS 读侧查询
Akka Persistence Query 实战指南:用统一异步流接口构建 CQRS 读侧查询 Akka Persistence Query 是 Akka 持
后端并发编程异步编程VictoriaMetrics 查询执行统计(Query Stats):慢查询日志与性能分析实战指南
VictoriaMetrics 查询执行统计(Query Stats):慢查询日志与性能分析实战指南 查询执行统计是 VictoriaMetrics 提供的查询
时序数据库数据库指标监控可观测性后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考