Clojure 领域建模与事件溯源:DDD、CQRS 与函数式实现

用 Clojure 实践 DDD 与事件溯源:限界上下文与聚合、decide/evolve 纯函数聚合建模、事件流与事件存储实现、CQRS 命令查询分离、投影与读模型、与 spec 校验集成、事件版本化与 Kafka/PostgreSQL 生产化,附完整可运行的银行账户事件溯源示例。

事件溯源(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 事件是唯一事实

事件溯源的三条铁律:

  1. 事件不可变:一旦发布,永不修改、永不删除(矫正用「补偿事件」)。
  2. 事件是事实:OrderConfirmed 表示「已发生的事实」,而非「将要执行的操作」。
  3. 状态可重放:任何时刻的状态 = 从头重放所有事件。

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 最佳实践清单

  1. decide/evolve 保持纯函数:所有 IO 移到处理器与投影层。
  2. 事件是外部契约:版本化事件类型,用 migrate 兼容旧数据。
  3. 聚合内强一致、跨聚合最终一致:用 saga 处理跨聚合流程。
  4. 读模型大胆重建:投影损坏直接重放事件流。
  5. 乐观并发:(stream_id, version) 唯一约束是底线。
  6. 为每个不变量写属性测试:余额非负、状态机合法转移等。

8.2 总结

概念Clojure 实现
DDD 聚合decide(命令校验)+ evolve(事件应用)
事件溯源不可变事件 + append-only 存储 + 重放
CQRS命令侧事件流 + 查询侧投影读模型
一致性聚合强一致 + 读模型最终一致
审计事件流天然是完整审计日志

Clojure 的函数式建模让 DDD 与事件溯源落地为一对纯函数:decide 承载业务规则,evolve 承载状态转换。不可变数据结构保证了事件永不篡改,persistent 版本共享让重放高效。这套架构尤其适合金融、电商、审计敏感等领域——数据永不丢失、逻辑可测试、状态可回放。进一步可结合 Clojure spec 与测试 强化命令校验,或参考 Kafka 专题 构建跨服务事件总线。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「clojure」更多文章

  1. Clojure 性能优化与 GraalVM 原生编译
  2. Clojure 生成式测试实战:test.check、收缩与 spec 集成
  3. ClojureScript 全栈开发:reagent、re-frame 与 shadow-cljs 实战