事件溯源(Event Sourcing)把状态重新定义为事件的累积:系统的真相不是当前数据库行,而是一串不可变的事实。配合 CQRS(命令查询职责分离)与 DDD(领域驱动设计),这种架构让业务逻辑变成纯函数、让审计成为第一公民、让时间旅行调试成为可能。Clojure 的不可变数据与函数式风格,与事件溯源几乎是天作之合。本文将从 DDD 聚合、事件溯源核心、CQRS 拆分,到完整的 Clojure 函数式实现与生产化考量,逐步展开。
后端数据持久化部分可参考 Clojure 数据库访问实战 与 PostgreSQL 专题;异步事件分发可参考 Clojure 并发设计模式。
1. DDD:领域驱动设计基础
1.1 限界上下文与通用语言
DDD 的第一步是划分限界上下文(Bounded Context)——每个上下文有独立的通用语言(Ubiquitous Language)与模型:
┌────────────────────────────────────────────────┐
│ 限界上下文:订单管理 │
│ 语言:Order, LineItem, Amount, Shipment │
│ ┌──────────────────┐ ┌──────────────────┐ │
│ │ 订单聚合 │ │ 库存聚合 │ │
│ │ Order 状态机 │ │ Stock 扣减 │ │
│ └──────────────────┘ └──────────────────┘ │
└────────────────────────────────────────────────┘
一个「Order」在订单上下文和物流上下文中的含义不同——上下文边界防住了模型膨胀。
1.2 聚合与聚合根
聚合(Aggregate)是一组作为一个一致性单位的领域对象,聚合根(Aggregate Root)是外部访问的唯一入口:
| 概念 | 说明 |
|---|---|
| 聚合根 | 如 Order,外部命令只能作用于根 |
| 聚合内部 | LineItem、Address 等不可直接从外部修改 |
| 一致性边界 | 聚合内所有约束在同一事务内保证 |
| 唯一标识 | 用 UUID 而非自增 ID,便于事件溯源 |
1.3 传统 CRUD vs 事件溯源
;; 传统 CRUD:状态被原地覆盖,历史丢失
(def order {:id "o1" :status :pending :items []})
(def order' (assoc order :status :confirmed)) ; 之前的状态被丢弃
;; 事件溯源:状态是事件的 fold,历史永不丢失
(def events [{:type ::order-created :id "o1" :items []}
{:type ::order-confirmed :id "o1"}])
(def current-state
(reduce apply-event {} events)) ; 从事件重放状态
2. 事件溯源核心
2.1 事件是唯一事实
事件溯源的三条铁律:
- 事件不可变:一旦发布,永不修改、永不删除(矫正用「补偿事件」)。
- 事件是事实:
OrderConfirmed表示「已发生的事实」,而非「将要执行的操作」。 - 状态可重放:任何时刻的状态 = 从头重放所有事件。
2.2 事件流与事件存储
;; 一个领域事件的标准形状
{:event-id #uuid"f81d4fae-7dec-11d0-a765-00a0c91e6bf6"
:stream-id "account-1001" ; 属于哪个聚合实例
:type :account/credited ; 事件类型(版本化见 6.2)
:timestamp #inst"2026-09-26T09:00:00.000Z"
:data {:amount 500 :ref "tx-1"}} ; 事件载荷(不可变 map)
事件存储(Event Store)按 stream-id 顺序追加。聚合实例的所有事件构成一条事件流:
account-1001: [AccountOpened] → [Credited 500] → [Debited 200] → [Credited 100]
2.3 投影与读模型
投影(Projection) 把事件流转成针对查询优化的读模型:
;; 从事件流投影出账户余额
(defn project-balance [events]
(reduce (fn [bal ev]
(case (:type ev)
:account/opened 0
:account/credited (+ bal (:amount (:data ev)))
:account/debited (- bal (:amount (:data ev)))
bal))
0
events))
3. CQRS:命令与查询分离
3.1 为什么拆分
| 命令侧(写) | 查询侧(读) | |
|---|---|---|
| 模型 | 聚合 + 事件 | 投影 + 物化视图 |
| 输入 | 命令(动词) | 查询(名词) |
| 目标 | 一致性、业务规则 | 性能、灵活性 |
| 存储 | 事件流 | 读模型表 |
| 一致性 | 强一致(聚合内) | 最终一致 |
单一 CRUD 模型往往「读写互相迁就」:读需求要反规范化,写需求要约束完整性。CQRS 让两侧各取所需。
3.2 命令处理流程
;; 命令处理器(Command Handler)五步:
;; 1. 加载聚合事件流
;; 2. 重放为当前状态
;; 3. 业务规则校验(decide)
;; 4. 生成新事件(decide 返回值)
;; 5. 追加事件并发布
(defn handle-command! [store {:keys [stream-id] :as command}]
(let [events (load-events store stream-id)
state (evolve events)
new-evs (decide state command)] ; 纯函数:校验 + 产出事件
(append-events! store stream-id new-evs)
(publish! new-evs)
new-evs))
4. Clojure 函数式建模
4.1 decide / evolve:纯函数聚合
事件溯源的核心模式是一对纯函数:
;; decide:命令 + 状态 → 新事件(校验业务规则,失败抛异常或返回错误)
;; 无副作用,可完美单元测试
;; evolve:状态 + 事件 → 新状态(纯 fold)
;; 幂等,可用于重放与投影
银行账户的完整聚合:
(ns bank.account
(:require [clojure.spec.alpha :as s]))
;; ---------- 事件 ----------
(s/def ::amount (s/and pos? number?))
(s/def ::event (s/keys :req-un [::type ::stream-id ::data]))
;; ---------- 状态 ----------
(defn initial-state [] {:balance 0 :version 0})
;; ---------- evolve:事件 → 状态 ----------
(defn evolve
([events] (reduce evolve (initial-state) events))
([state {:keys [type data]}]
(case type
:account/opened (-> state (assoc :balance 0) (update :version inc))
:account/credited (-> state (update :balance + (:amount data)) (update :version inc))
:account/debited (-> state (update :balance - (:amount data)) (update :version inc))
:account/closed (assoc state :closed? true :version (inc (:version state)))
state)))
;; ---------- decide:命令 → 事件(业务规则所在地) ----------
(defn decide
[{:keys [balance closed?]} {:keys [type amount] :as cmd}]
(case type
:open-account
[{:type :account/opened :data {}}]
:credit-account
[{:type :account/credited :data {:amount amount :ref (:ref cmd)}}]
:debit-account
(cond
closed? (throw (ex-info "账户已关闭" {:cmd cmd}))
(< balance amount) (throw (ex-info "余额不足" {:cmd cmd :balance balance}))
:else [{:type :account/debited :data {:amount amount :ref (:ref cmd)}}])
:close-account
(if (pos? balance)
(throw (ex-info "请先清零余额" {:balance balance}))
[{:type :account/closed :data {}}])))
decide 与 evolve 分离后,业务规则(decide)不依赖任何 IO,全部可用纯函数属性测试(参考 Clojure 生成式测试实战)。
4.2 用 map 与 record 表示领域对象
;; 领域对象优先用不可变 map(简单、可序列化)
(def account {:id "acc-1001" :balance 500})
;; 需要行为与协议时用 defrecord
(defrecord Account [id balance closed?]
AccountAPI
(debit [acc amount] (assoc acc :balance (- balance amount))))
事件溯源风格下更推荐纯 map + 纯函数:AccountAPI 协议提供 decide/evolve 的分派入口,具体行为仍是纯函数。
4.3 事件存储实现(PostgreSQL + next.jdbc)
(require '[next.jdbc :as jdbc])
;; 建表
(def create-event-store-sql
["CREATE TABLE IF NOT EXISTS event_store (
id BIGSERIAL PRIMARY KEY,
stream_id TEXT NOT NULL,
version BIGINT NOT NULL,
event_type TEXT NOT NULL,
payload JSONB NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
UNIQUE (stream_id, version) -- 并发防重:乐观锁
)"])
(defn load-events [ds stream-id]
(jdbc/execute! ds
["SELECT payload, version FROM event_store
WHERE stream_id = ? ORDER BY version"
stream-id]
{:builder-fn next.jdbc.result-set/as-unqualified-maps}))
(defn append-events! [ds stream-id expected-version events]
;; expected-version 实现乐观并发控制
(doseq [[i ev] (map-indexed vector events)]
(let [v (+ expected-version i 1)]
(try
(jdbc/execute-one! ds
["INSERT INTO event_store (stream_id, version, event_type, payload)
VALUES (?, ?, ?, ?::jsonb)"
stream-id v (name (:type ev))
(json/write-str (:data ev))])
(catch java.sql.SQLIntegrityConstraintViolationException e
(throw (ex-info "并发冲突,请重试" {:stream-id stream-id})))))))
UNIQUE (stream_id, version) 是事件溯源的并发防线:两个命令同时写同一版本时,一个成功、一个违反唯一约束重试。
5. CQRS 读侧:投影与查询
5.1 读模型表
;; 投影到 accounts 读模型(balance 快照)
(defn project-to-account! [ds ev]
(let [acc (jdbc/get-by-id ds :accounts (:stream-id ev))]
(if acc
(jdbc/update! ds :accounts
{:balance (apply-event-balance (:balance acc) ev)}
{:id (:stream-id ev)})
(jdbc/insert! ds :accounts
{:id (:stream-id ev) :balance (event-balance ev)}))))
;; 订阅事件流,增量更新读模型
(defn projection-loop [ds events-ch]
(go-loop []
(when-let [ev (<! events-ch)]
(project-to-account! ds ev)
(recur))))
5.2 读模型 vs 事件流
| 特性 | 事件流(写侧) | 读模型(读侧) |
|---|---|---|
| 存储 | append-only | 可覆盖(从事件重建) |
| 一致性 | 强一致 | 最终一致 |
| 重建 | 不可重建 | 可随时从事件流重建 |
| 查询能力 | 弱(需重放) | 强(按需建模) |
读模型可以随意丢弃重建——这是 CQRS 的容错红利:任何读模型损坏,重放事件流即可恢复。
5.3 查询处理器
(defn get-account-summary [ds account-id]
(jdbc/get-by-id ds :accounts account-id))
(defn get-account-statement [ds account-id]
(jdbc/execute! ds
["SELECT payload, created_at FROM event_store
WHERE stream_id = ? ORDER BY version" account-id]))
6. 生产化考量
6.1 与消息队列集成
事件溯源常与 Kafka/消息队列配合:事件存储追加后发布到 bus,投影与外部系统订阅:
;; 事件发布(示例伪代码,真实实现见 Kafka 专题)
(defn publish! [events]
(doseq [ev events]
(kafka/produce! "domain-events" (json/write-str ev))))
;; 投影消费者
(defn projection-consumer [ds]
(go-loop []
(let [ev (<! (:events-ch bus))]
(project-to-account! ds ev)
(recur))))
异步事件分发细节可参考 Clojure 并发设计模式 的管道模式与 Kafka 专题。
6.2 事件版本化
领域演进时事件结构会变。事件类型要带版本,迁移策略有二:
;; 策略一:事件类型带版本号
{:type :account/credited-v2
:data {:amount 500 :currency "CNY"}} ; v1 只有 amount
;; 策略二:统一 migrate 函数
(defn migrate-event [{:keys [type data] :as ev}]
(case type
:account/credited (assoc ev :data (assoc data :currency "CNY"))
ev))
| 策略 | 适用 |
|---|---|
| 类型带版本 | 事件形态变化大 |
| 统一 migrate | 只需兼容旧事件 |
6.3 一致性模型选择
| 场景 | 一致性 |
|---|---|
| 余额扣减(聚合内) | 强一致(事件追加 + 唯一约束) |
| 读模型投影 | 最终一致(异步) |
| 跨聚合事务 | 事件驱动 + saga 补偿 |
| 报表/分析 | 最终一致(批量重放) |
6.4 常见陷阱
| 陷阱 | 后果 | 对策 |
|---|---|---|
| 事件载荷耦合内部结构 | 重构困难 | 事件是外部契约,用版本化 + 迁移 |
| decide 里做 IO | 不可测、难重放 | 保持纯函数,IO 只在处理器层 |
| 并发无保护 | 事件乱序、丢事件 | UNIQUE(stream_id, version) 乐观锁 |
| 大聚合(过宽状态) | 重放慢、锁竞争 | 拆聚合、按子流建模 |
| 投影与写侧混用 | CQRS 边界被打破 | 读侧一律走读模型 |
7. 完整实战:银行转账系统
(ns bank.system
(:require [bank.account :as acct]
[next.jdbc :as jdbc]
[clojure.spec.alpha :as s]
[clojure.spec.test.alpha :as stest]))
;; ---------- 命令校验(spec) ----------
(s/def ::command (s/keys :req-un [::acct/stream-id ::acct/type ::acct/data]))
;; ---------- 应用服务(编排层) ----------
(defn handle [ds {:keys [stream-id] :as cmd}]
(s/assert ::command cmd)
(let [events (load-events ds stream-id)
state (acct/evolve events)]
(->> (acct/decide state cmd)
(append-events! ds stream-id (:version state))
(run! publish!))))
;; ---------- 使用 ----------
(def ds (jdbc/get-datasource {:dbtype "postgresql" :dbname "bank"}))
;; 开户
(handle ds {:stream-id "acc-1" :type :open-account :data {}})
;; 存款 500
(handle ds {:stream-id "acc-1" :type :credit-account :data {:amount 500}})
;; 取款 200
(handle ds {:stream-id "acc-1" :type :debit-account :data {:amount 200}})
;; 取款超余额 → 抛出 "余额不足"
(handle ds {:stream-id "acc-1" :type :debit-account :data {:amount 999}})
;; ---------- 审计与时间旅行 ----------
(get-account-statement ds "acc-1")
;; => 完整事件流:开户、存 500、取 200、失败的取款尝试(如果记录了)
;; 在任何时刻重建状态
(->> (load-events ds "acc-1") acct/evolve)
;; => {:balance 300, :version 3}
7.1 用属性测试守护 decide/evolve
(defspec debit-never-goes-negative
100
(prop/for-all [start gen/pos-int
amount gen/pos-int]
(let [events [{:type :account/opened :data {}}
{:type :account/credited :data {:amount start}}]
state (acct/evolve events)]
(if (>= start amount)
(>= (:balance (acct/evolve (conj events
{:type :account/debited :data {:amount amount}}))) 0)
(thrown? clojure.lang.ExceptionInfo
(acct/decide state {:type :debit-account :amount amount}))))))
这条属性测试保证任何随机取款序列都不会让余额变负——正是账户系统的核心不变量。
8. 最佳实践与总结
8.1 最佳实践清单
- decide/evolve 保持纯函数:所有 IO 移到处理器与投影层。
- 事件是外部契约:版本化事件类型,用 migrate 兼容旧数据。
- 聚合内强一致、跨聚合最终一致:用 saga 处理跨聚合流程。
- 读模型大胆重建:投影损坏直接重放事件流。
- 乐观并发:
(stream_id, version)唯一约束是底线。 - 为每个不变量写属性测试:余额非负、状态机合法转移等。
8.2 总结
| 概念 | Clojure 实现 |
|---|---|
| DDD 聚合 | decide(命令校验)+ evolve(事件应用) |
| 事件溯源 | 不可变事件 + append-only 存储 + 重放 |
| CQRS | 命令侧事件流 + 查询侧投影读模型 |
| 一致性 | 聚合强一致 + 读模型最终一致 |
| 审计 | 事件流天然是完整审计日志 |
Clojure 的函数式建模让 DDD 与事件溯源落地为一对纯函数:decide 承载业务规则,evolve 承载状态转换。不可变数据结构保证了事件永不篡改,persistent 版本共享让重放高效。这套架构尤其适合金融、电商、审计敏感等领域——数据永不丢失、逻辑可测试、状态可回放。进一步可结合 Clojure spec 与测试 强化命令校验,或参考 Kafka 专题 构建跨服务事件总线。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。