持久化与事件溯源:journal、快照与 CQRS 读模型

事件溯源把系统状态定义为事件序列的折叠结果:写模型只追加不可变事件,读模型由投影派生,命令查询职责分离因此成为它的自然搭档。本文讲清持久化 actor 的命令与事件处理模型、日志与快照如何协作、投影与读模型的构建方式、事件版本演进与兼容策略,并给出恢复性能调优与日常运维的实践要点。

「把当前状态存进数据库」是最自然的做法,也是最难审计的做法:一旦业务问「这个订单金额是怎么变成 380 的」,行级快照给不出答案。事件溯源(Event Sourcing)换了个建模方向——状态不是被存储的,而是被推导出来的:系统只追加不可变事件,当前状态是事件序列折叠的结果。任何时刻的状态都可重建、可回放、可审计。

代价是复杂度:你要处理事件 schema 演进、快照与恢复性能、读模型的最终一致性、投影的幂等与重放。Akka/Pekko Persistence 把这些机制做成了可复用的基础设施。本文按「写模型 → 持久化 → 读模型 → 演进 → 运维」的顺序,把这些机制讲透。

前置:Actor 模型详解 、Akka Cluster 与分片 。

1. 事件溯源与 CQRS 的模型

事件溯源的核心约定有三条:事件是事实、只追加不修改;状态由事件折叠得出;副作用只发生在事件被持久化之后。命令(Command)表达意图,可能被拒绝;事件(Event)表达已发生的事实,一旦写入不可撤销。

模型分层:
Command  意图,可拒绝(余额不足 → 拒绝)
  ↓ 校验 + 决策
Event    事实,不可变,只追加(Deposited、Withdrawn)
  ↓ 折叠
State    由事件序列推导出的当前状态(balance = sum(deposits) - sum(withdrawals))
  ↓ 投影
ReadModel 面向查询的形态(账户余额表、交易流水视图)

CQRS 的职责分离:
写侧:命令 → 校验 → 事件 → journal(强一致、单实体串行)
读侧:journal → 投影 → 读模型(最终一致、可任意重塑)
import org.apache.pekko.persistence.typed.scaladsl.*
import org.apache.pekko.persistence.typed.PersistenceId

final case class Account(balance: BigDecimal, currency: String)

sealed trait Cmd
final case class Deposit(amount: BigDecimal) extends Cmd
final case class Withdraw(amount: BigDecimal) extends Cmd

sealed trait Evt
final case class Deposited(amount: BigDecimal) extends Evt
final case class Withdrawn(amount: BigDecimal) extends Evt

val commandHandler: (Account, Cmd) => Effect[Evt, Account] =
  (state, cmd) => cmd match
    case Deposit(a)  => Effect.persist(Deposited(a))
    case Withdraw(a) if state.balance >= a => Effect.persist(Withdrawn(a))
    case Withdraw(_) => Effect.none          // 拒绝:不产生事件,状态不变

val eventHandler: (Account, Evt) => Account =
  (state, evt) => evt match
    case Deposited(a) => state.copy(balance = state.balance + a)
    case Withdrawn(a) => state.copy(balance = state.balance - a)

工程要点:命令处理函数必须是纯的——它只能返回「要持久化哪些事件」,绝不能在里面写数据库、发消息、打外部 API。所有副作用应当挂在 Effect.thenRun(...) 或 Effect.persist(...).thenRun(...) 上,这样在事件确认写入之前不会有任何外部可见的副作用。

2. 持久化 actor 与恢复

EventSourcedBehavior 把「命令处理」「事件应用」「初始状态」三段组合成一个持久化 actor。它的生命周期里有一段特殊的恢复(recovery)阶段:actor 启动时先从 journal 回放历史事件重建状态,回放完成前不接收任何命令。

import org.apache.pekko.actor.typed.{ActorSystem, Behavior}
import org.apache.pekko.actor.typed.scaladsl.Behaviors
import org.apache.pekko.persistence.typed.PersistenceId
import org.apache.pekko.persistence.typed.scaladsl.EventSourcedBehavior

object AccountEntity:
  def apply(id: String): Behavior[Cmd] =
    EventSourcedBehavior[Cmd, Evt, Account](
      persistenceId = PersistenceId.ofUniqueId(s"account-$id"),
      emptyState = Account(BigDecimal(0), "CNY"),
      commandHandler = commandHandler,
      eventHandler = eventHandler
    ).withTagger(_ => Set("account"))       // 打 tag,供投影按流消费
恢复流程(每个持久化 actor 启动时):
1. 读最新快照(若有)→ 得到基线状态与序号 seqNr
2. 从 seqNr+1 起回放 journal 事件,逐个 apply
3. 回放到末尾 → 状态就绪 → 开始接收命令
4. 回放期间收到的命令被暂存(stash),就绪后按序处理
恢复的代价:
□ 事件越多恢复越慢(与事件数线性相关)
□ 快照把恢复起点前移,代价是快照本身的序列化与存储
□ 大事件(如整张表快照)会让回放与序列化都变慢
// 监听恢复完成:常用于「恢复后再做外部动作」
EventSourcedBehavior[Cmd, Evt, Account](...)
  .receiveSignal {
    case (state, RecoveryCompleted) =>
      // 此处可安全读取外部系统、上报指标
      ()
    case (_, RecoveryFailed(cause)) =>
      // 恢复失败通常意味着 journal 不可用或事件不兼容
      throw cause
  }

工程要点:恢复阶段不接收命令这一点决定了可用性——一个有一百万事件、没有快照的实体,每次重启都要回放一百万条。把「快照间隔」和「单实体事件数上限」一起设计:如果某个实体的预期事件数会长期增长,就该考虑把它拆成多个短生命周期实体(例如按月份分片)。

3. Journal 与快照的协作

Journal 是事件的权威存储(append-only),快照(Snapshot)是性能优化而非真相来源——它随时可以被删除重建,因为事件永远在。两者配合的关键参数是快照触发条件与保留策略。

# application.conf:journal 与快照的配置
pekko.persistence.journal.plugin = "pekko.persistence.journal.jdbc"
pekko.persistence.snapshot-store.plugin = "pekko.persistence.snapshot-store.jdbc"

# 快照策略:每 N 个事件或每 N 条事件后触发
pekko.persistence.journal.jdbc {
  class = "org.apache.pekko.persistence.jdbc.journal.JdbcAsyncWriteJournal"
  tables {
    event_journal { tableName = "event_journal" }
    event_tag     { tableName = "event_tag" }
  }
}
pekko.persistence.snapshot-store.jdbc {
  class = "org.apache.pekko.persistence.jdbc.snapshot.JdbcSnapshotStore"
  tables { snapshot { tableName = "snapshot" } }
}
// 快照策略:按事件数触发(例如每 100 个事件存一次)
import org.apache.pekko.persistence.typed.scaladsl.RetentionCriteria

EventSourcedBehavior[Cmd, Evt, Account](...)
  .withRetention(RetentionCriteria
    .snapshotEvery(numberOfEvents = 100, keepNSnapshots = 2))
  // 语义:每 100 个事件存一次快照,只保留最近 2 个快照
  // 注意:默认不会删除 journal 事件(事件是真相来源)
快照相关的关键参数与取舍:
□ snapshotEvery(n):n 越小恢复越快,但序列化与存储开销越大
□ keepNSnapshots(k):保留几份快照;k=1 意味着快照写失败就没退路
□ deleteEventsOnSnapshot:是否裁剪旧事件(默认 false,谨慎开启)
  → 开启后事件不再可回放,审计与投影重建能力受损
□ 快照序列化格式:与事件一样面临 schema 演进问题
□ 快照存储:与 journal 同库可省运维成本,但会争抢连接与 IO

工程要点:除非有明确的合规与容量约束,不要开启 deleteEventsOnSnapshot。裁剪事件会永久失去回放能力——投影重建、事件审计、新读模型的回填全部依赖事件。容量问题的正确解法是分片(按时间或按聚合根拆实体),而不是删事件。

4. 投影与读模型

读模型(Read Model)由投影(Projection)从事件流派生。Akka/Pekko 的 Projection API 支持按 tag 消费事件流、管理 offset、批量处理与失败重启。核心要求是幂等——投影可能因重启、rebalance、offset 回退而重复处理同一条事件。

import org.apache.pekko.projection.eventsourced.*
import org.apache.pekko.projection.eventsourced.scaladsl.EventSourcedProvider
import org.apache.pekko.projection.scaladsl.{Projection, SourceProvider}
import org.apache.pekko.projection.jdbc.scaladsl.JdbcProjection

// 按 tag 消费事件流(event_journal 的 event_tag 表)
val sourceProvider: SourceProvider[Offset, EventEnvelope[Evt]] =
  EventSourcedProvider.eventsByTag[Evt](system, readJournalPluginId = "pekko.persistence.jdbc.query", tag = "account")

val projection: Projection[EventEnvelope[Evt]] =
  JdbcProjection.exactlyOnce(
    projectionId = ProjectionId("account-balance-view", "0"),
    sourceProvider = sourceProvider,
    sessionFactory = () => sessionFactory,
    handler = () => new AccountBalanceHandler   // 在 handler 里写读模型
  )
import org.apache.pekko.projection.eventsourced.EventEnvelope

// 幂等投影:把 offset 与业务写入放进同一个事务
class AccountBalanceHandler extends JdbcHandler[EventEnvelope[Evt], Session]:
  def process(session: Session, envelope: EventEnvelope[Evt]): Unit =
    envelope.event match
      case Deposited(a) =>
        // upsert 而非 insert:重复投递不会产生重复行
        session.prepareStatement(
          """INSERT INTO account_balance(account_id, balance)
             |VALUES (?, ?)
             |ON CONFLICT (account_id) DO UPDATE
             |SET balance = account_balance.balance + EXCLUDED.balance""".stripMargin)
          .setString(1, envelope.persistenceId)
          .setBigDecimal(2, a.bigDecimal)
          .executeUpdate()
      case Withdrawn(a) => // 对称处理
投影的三种一致性级别:
□ at-least-once  最简单:可能重复,业务侧必须幂等
□ at-least-once + 幂等键  推荐:用事件序号或事件 id 做去重
□ exactly-once  需要:把 offset 提交与读模型写入放进同一事务
   (要求读模型与 offset 存储在同一数据库)
投影的运维要点:
□ 按 tag 分区:同一 tag 的事件必须按序处理(同一实体的顺序是硬要求)
□ offset 存储:与读模型同库才能做事务性 exactly-once
□ 重放:读模型结构变更时,重置 offset 从头重建
□ 背压:投影处理慢会累积 lag,需监控

工程要点:投影 handler 里的写入必须幂等,最可靠的做法是把「事件 id(或 persistenceId + sequenceNr)」作为一个唯一键写进读模型表,重复投递触发 ON CONFLICT DO NOTHING。这比依赖「offset 一定不重复」稳健得多——rebalance 与重启都会破坏这个假设。

5. 事件版本演进与兼容

事件一旦写入就永久存在,schema 演进必须向后兼容:新代码要能读旧事件,旧代码(回滚时)最好也能读新事件。Pekko 提供 EventAdapter 做「旧事件 → 新事件」的转换。

import org.apache.pekko.persistence.typed.*
import org.apache.pekko.persistence.typed.scaladsl.EventSourcedBehavior

// 场景:v1 的 Deposited(amount) 演进为 v2 的 Deposited(amount, currency, ts)
final case class DepositedV1(amount: BigDecimal)
final case class DepositedV2(amount: BigDecimal, currency: String, ts: Long)

// upcasting:把历史事件转成当前形态
class DepositedUpcaster extends EventAdapter[Any, Evt]:
  override def toJournal(e: Evt): Any = e
  override def manifest(event: Evt): String = event.getClass.getSimpleName
  override def fromJournal(p: Any, manifest: String): EventSeq[Evt] = p match
    case DepositedV1(a) => EventSeq.single(DepositedV2(a, "CNY", 0L))  // 补默认值
    case e: Evt         => EventSeq.single(e)
事件演进的安全操作:
✓ 新增字段 + 给默认值(upcaster 补齐旧事件)
✓ 新增事件类型(读侧有 unknown 兜底)
✓ 放宽字段约束(非空 → 可空)
✗ 删除字段(旧代码读到会失败)
✗ 改字段类型(BigDecimal → String 之类)
✗ 改事件类名/包名(manifest 失配)
✗ 改字段语义(同名不同义,最隐蔽)
演进策略:
□ manifest 显式化:不要用默认 FQCN,用稳定的字符串标识
□ 多版本共存:读侧同时认识 v1 与 v2,写侧只写 v2
□ 大爆炸改造:必要时新开一条事件流,老流冻结只读

工程要点:给每个事件类显式声明一个稳定的 manifest 字符串(如 order-deposited-v1)。默认用全限定类名作 manifest,一旦重构改了包名或类名,历史事件就再也读不出来——这是事件溯源里最昂贵的一类事故。

6. 快照与恢复性能调优

恢复时间是事件溯源系统最直接的可用性指标。优化路径有三条:快照前移恢复起点、降低单事件回放成本、控制单实体的事件总量。

恢复耗时模型(近似):
恢复时间 ≈ 读快照 + (总事件数 - 快照位置) × 单事件回放成本
优化杠杆:
1. snapshotEvery(n) 调小 → 减少回放事件数(代价:快照写入与存储)
2. 单事件回放成本:避免大对象、避免回放里做复杂计算
3. 事件总数:按时间/业务维度拆实体(如 account-2026-10)
4. 并行恢复:多实体并发恢复,受 journal 读带宽限制
5. 恢复限流:pekko.persistence.journal.jdbc 的读并发配置
// 把「实体生命周期」设计成有界:按周期分片
def entityId(account: String, period: String): String = s"$account-$period"

// 读侧需要跨周期聚合时,用投影做「周期实体 → 全局读模型」的合并
// 这样单实体的事件数天然有界,恢复时间可预测
val pid = PersistenceId.ofUniqueId(entityId("A100", "2026-10"))
监控指标:
□ 恢复耗时 p50/p99(按实体类型分组)
□ 单实体事件数分布(长尾实体是恢复慢的元凶)
□ 快照命中率与快照大小
□ journal 写入延迟(影响命令处理的响应时间)
□ 投影 lag(读模型落后事件流多少条)

工程要点:单实体事件数应当有设计上限。如果一个聚合根的生命周期是永久的(账户、租户),它的事件会无限增长,恢复时间也随之线性增长。把聚合根按时间或业务周期切分(account-2026-10),让每个实体的事件数有界,是唯一可持续的解法。

7. 集群分片与持久化

把持久化实体放到集群里跑,需要 Cluster Sharding 保证「同一实体同一时刻只有一个实例」。分片与持久化组合时,实体在节点间迁移会触发先停后启:旧节点停止实体,新节点恢复实体。

import org.apache.pekko.cluster.sharding.typed.scaladsl.{ClusterSharding, Entity}
import org.apache.pekko.cluster.sharding.typed.ShardingMessageExtractor
import org.apache.pekko.actor.typed.scaladsl.Behaviors

val sharding = ClusterSharding(system)
val region = sharding.init(
  Entity(EntityTypeKey[Cmd]("account")) { ctx =>
    AccountEntity(ctx.entityId)
  }.withMessageExtractor(extractor)
)
// 发消息:按实体 id 路由,运行时决定落在哪个节点
region ! ShardingEnvelope("account-A100", Deposit(BigDecimal(100)))
分片与持久化的交互要点:
□ 实体迁移 = 停止 + 恢复,恢复期间该实体不可用(有短暂不可用窗口)
□ 迁移频率由 rebalance 策略与节点稳定性决定
□ 恢复慢的实体会拖长迁移窗口,放大不可用时间
□ 分片数量固定后不可减;实体数远大于分片数才能均衡
□ 与 projection 配合:投影按 tag 消费,与分片位置无关
□ 灰度升级:滚动重启会让实体多次迁移,恢复性能是升级速度的上限

工程要点:恢复性能直接决定滚动升级的速度上限。Kubernetes 滚动更新时,节点上的实体会迁移到新节点并逐个恢复;如果单实体恢复要几秒,一次升级就可能造成大面积不可用。上线前务必在预发环境做一次全量滚动重启,测量最慢实体的恢复时间。

8. 测试与运维工具

事件溯源的测试优势在于确定性:给定同一组事件,折叠出的状态必然相同。这让「事件回放测试」成为最有效的回归手段。

import org.apache.pekko.persistence.testkit.scaladsl.EventSourcedBehaviorTestKit
import org.apache.pekko.actor.testkit.typed.scaladsl.ScalaTestWithActorTestKit
import org.scalatest.wordspec.AnyWordSpecLike

class AccountSpec extends ScalaTestWithActorTestKit(AccountSpec.config)
    with AnyWordSpecLike:

  private val kit = EventSourcedBehaviorTestKit[AccountEntity.Cmd, Evt, Account](system, AccountEntity("A1"))

  "Account" should {
    "接受合法取款" in {
      kit.runCommand(Deposit(BigDecimal(100)))
      val result = kit.runCommand(Withdraw(BigDecimal(30)))
      assert(result.event == Withdrawn(BigDecimal(30)))
      assert(result.state.balance == BigDecimal(70))
    }
    "拒绝超额取款" in {
      val result = kit.runCommand(Withdraw(BigDecimal(999)))
      assert(result.hasNoEvents)          // 拒绝 = 不产生事件
      assert(result.state.balance == BigDecimal(0))
    }
  }
运维工具与流程:
□ 事件回放:从 journal 重放到指定序号,验证状态推导
□ 投影重建:重置 offset,从头投影出新的读模型
□ 事件审计:按实体查询事件序列(谁在什么时候做了什么)
□ 一致性检查:对比「事件折叠的状态」与「读模型的当前值」
□ 灰度:新版本先只读(投影)后写(命令),观察事件兼容性
□ 告警:恢复耗时、投影 lag、journal 写延迟、快照失败率

工程要点:定期跑一次一致性检查——对抽样实体,分别用「折叠事件」与「查询读模型」两条路径算状态,比对是否一致。这是发现投影 bug、offset 跳号、幂等失效的最直接手段,比任何日志都有效。

9. 速查表与一句话记忆

问题一句话答案
事件与状态的关系状态是事件折叠的结果,事件是唯一真相
命令能被拒绝吗能,拒绝即不产生事件、状态不变
命令处理函数能写副作用吗不能,副作用挂在 Effect.thenRun 上
快照是什么恢复性能优化,可随时重建,不是真相来源
快照间隔怎么定由「可接受的恢复时间 ÷ 单事件回放成本」反推
事件能删吗除合规要求外不要删,删了就失去回放能力
投影为什么要幂等重启与 rebalance 会重复投递同一事件
exactly-once 的前提offset 与读模型写入在同一事务、同一数据库
演进的安全区新增字段给默认值 + upcaster 补齐旧事件
manifest 用什么稳定的自定义字符串,绝不用默认 FQCN
恢复慢怎么办快照前移起点 + 拆实体让单实体事件数有界

一句话记忆:事件溯源 = 事件是唯一真相(只追加)+ 状态是折叠结果(可重建)+ 快照只做恢复优化(可丢弃)+ 读模型靠投影派生(最终一致、必须幂等)+ 演进靠 upcaster 与稳定 manifest(向后兼容)+ 恢复性能靠快照间隔与实体切分(事件数有界)——写侧单实体串行保证强一致,读侧任意重塑支撑查询演化。

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「scala」更多文章

  1. Scala CLI 与工具链现代化:指令声明、打包发布与 CI 集成
  2. JSON 编解码与性能:circe 派生、jsoniter-scala 代码生成与选型
  3. fs2 流处理实战:Stream、Pipe 与背压模型