后端并发编程异步编程【免费下载链接】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-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。sbtlibraryDependencies com.typesafe.akka %% akka-persistence-query % AkkaVersionMavenproperties akka.version2.9.x/akka.version scala.binary.version2.13/scala.binary.version /properties dependency groupIdcom.typesafe.akka/groupId artifactIdakka-persistence-query_${scala.binary.version}/artifactId version${akka.version}/version /dependencyGradledef 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.leveldbJava 移植版或org.fusesource.leveldbjniJNI 原生版。3. 如何获取 ReadJournalReadJournal通过akka.persistence.query.PersistenceQuery扩展extension获取。核心要点使用LeveldbReadJournal.Identifier常量值为akka.persistence.query.journal.leveldb作为插件标识该标识同时是配置文件中的绝对路径二者一一对应。Scala完整代码见 LeveldbPersistenceQueryDocSpec.scalaimport akka.persistence.query.PersistenceQuery import akka.persistence.query.journal.leveldb.scaladsl.LeveldbReadJournal val queries PersistenceQuery(system).readJournalForLeveldbReadJournalJava完整代码见 LeveldbPersistenceQueryDocTest.javaimport 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 与 CurrentEventsByPersistenceIdeventsByPersistenceId用于检索某个指定persistenceId的PersistentActor所持久化的全部事件。Scalaval 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)JavaLeveldbReadJournal queries PersistenceQuery.get(system) .getReadJournalFor(LeveldbReadJournal.class, LeveldbReadJournal.Identifier()); SourceEventEnvelope, NotUsed source queries.eventsByPersistenceId(some-persistence-id, 0, Long.MAX_VALUE);语义要点文档明确约定源码注释亦有完整复述区间检索可通过fromSequenceNr与toSequenceNr限定事件子集使用0L和Long.MaxValueJava 为Long.MAX_VALUE则可检索全部事件。每个事件的序号会体现在EventEnvelope中因此可以从某个序号之后继续实现断点续流。顺序保证返回流按序号sequence number排序即与PersistentActor持久化事件的顺序一致。同一查询多次执行返回相同前缀、相同顺序的流元素除非事件被删除。Live 与 Current 的区别eventsByPersistenceId是活流——到达当前已存事件末尾时不会结束而是持续推送新持久化的事件currentEventsByPersistenceId则是当前快照流——到达已存事件末尾即完成。批处理刷新LevelDB 写日志会在事件持久化后尽快通知查询端但出于效率考虑查询端按批次拉取事件最多可能延迟到配置的refresh-interval或显式传入的RefreshInterval提示时长。失败语义若后端日志执行查询失败流将以 failure 结束。4.2 PersistenceIds 与 CurrentPersistenceIdspersistenceIds用于检索所有持久化 Actor 的persistenceId列表。Scalaval queries PersistenceQuery(system).readJournalForLeveldbReadJournal val src: Source[String, NotUsed] queries.persistenceIds()JavaLeveldbReadJournal queries PersistenceQuery.get(system) .getReadJournalFor(LeveldbReadJournal.class, LeveldbReadJournal.Identifier()); SourceString, NotUsed source queries.persistenceIds();语义要点无序返回流不保证顺序多次执行可能得到不同的排列顺序。Live 与 Current 的区别persistenceIds是活流新持久化 Actor 出现时会持续推送新的persistenceIdcurrentPersistenceIds则在遍历完当前集合后完成。无轮询与其他两种查询不同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 与 CurrentEventsByTageventsByTag用于检索带有指定 tag 的事件典型场景是某个聚合根Aggregate Root类型的所有领域事件。Scalaval queries PersistenceQuery(system).readJournalForLeveldbReadJournal val src: Source[EventEnvelope, NotUsed] queries.eventsByTag(tag green, offset Sequence(0L))JavaLeveldbReadJournal queries PersistenceQuery.get(system) .getReadJournalFor(LeveldbReadJournal.class, LeveldbReadJournal.Identifier()); SourceEventEnvelope, NotUsed source queries.eventsByTag(green, new Sequence(0L));如何给事件打标签WriteEventAdapter Tagged要给事件打上 tag需要创建一个 Event Adapters写事件适配器把事件包装进akka.persistence.journal.Tagged并附上tags集合。ScalaMyTaggingEventAdapter摘自 LeveldbPersistenceQueryDocSpec.scalaimport 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.javaimport 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; SetString tags new HashSetString(); 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 类型直接抛出IllegalArgumentExceptionLevelDB does not support ... offsets——即 LevelDB 实现不支持TimeBasedUUID等其他 offset。Offset 是排他的与 offset 序号完全相同的那条事件不会被包含在返回流中。这意味着你可以把EventEnvelope中返回的 offset 直接作为下一次查询的offset参数实现精确续传、不重不漏。Envelope 附加信息除 offset 外EventEnvelope还提供persistenceId与sequenceNr。其中sequenceNr是持久化该事件的 Actor 自己的序号persistenceIdsequenceNr构成事件的唯一标识。排序与稳定性返回流按 offsettag 序号排序与写日志存储顺序一致多次执行返回相同元素、相同顺序。删除不影响 tag 流文档用专门的 note 强调——通过deleteMessages(toSequenceNr)删除的事件不会从 tagged streamtag 事件流中删除。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 }各配置项说明配置项默认值作用classakka.persistence.query.journal.leveldb.LeveldbReadJournalProviderReadJournal 的 Provider 实现类一般无需修改write-plugin空字符串所连接的写日志插件配置的绝对路径必须是LeveldbJournal或SharedLeveldbJournal留空则使用akka.persistence.journal.plugin指定的默认日志refresh-interval3s写日志在事件持久化后立即通知查询端但查询端按批次拉取最多可延迟到该时长作用于eventsByPersistenceId与eventsByTagmax-buffer-size100单次查询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, sLeveldb 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否则直接抛出IllegalArgumentExceptionLeveldb read journal can only work with a Leveldb write journal。这是读日志与写日志必须配套的强约束——即使write-plugin配的是SharedLeveldbJournal读端校验依旧按此执行。refreshInterval只在 live 查询中传入 Stagecurrent*查询传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-sharedSharedLeveldbJournal。测试环境通常这样配置摘自 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.mdLevelDB 写日志实现见 LeveldbJournal.scala 及其配套的 LeveldbStore.scala。赞分享后端并发编程异步编程【免费下载链接】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 完全指南基于 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),仅供参考