当业务发展到一定规模,传统 CRUD 读写同库模型在复杂查询/审计/并发控制遇到瓶颈。CQRS + ES 是打破瓶颈的高级模式,但复杂度成倍上升,必须慎重选型。
Event Sourcing:事件为真相源
不存当前状态,只存导致状态变化的事件流(OrderCreated → ItemAdded → OrderPlaced)。重放所有事件可重建任意时间点状态,天然拥有完整审计轨迹。配合快照避免每次从 0 重放。
from __future__ import annotations
import json, time
from dataclasses import dataclass, field
from typing import Callable, Iterable, Any
from collections import defaultdict
@dataclass
class DomainEvent:
event_id: str
aggregate_type: str
aggregate_id: str
event_type: str
payload: dict
occurred_at: float = field(default_factory=time.time)
version: int = 1
class EventStore:
def __init__(self) -> None:
self._events: list[DomainEvent] = []
self._snapshots: dict[tuple[str, str], tuple[int, Any]] = {}
self._subs: dict[str, list[Callable[[DomainEvent], None]]] = defaultdict(list)
def append(self, events: Iterable[DomainEvent], expected_version: int | None = None) -> None:
events = list(events)
if expected_version is not None and events:
at, aid = events[0].aggregate_type, events[0].aggregate_id
cur = sum(1 for e in self._events if e.aggregate_type == at and e.aggregate_id == aid)
if cur != expected_version:
raise ValueError(f"Concurrency conflict: expected v{expected_version}, got v{cur}")
for e in events:
self._events.append(e)
for s in self._subs.get(e.event_type, []) + self._subs.get("*", []): s(e)
def load(self, agg_type: str, agg_id: str, snapshot_every: int = 50) -> tuple[Any, int]:
key = (agg_type, agg_id)
start_ver, state = self._snapshots.get(key, (0, None))
version = start_ver
events = [e for e in self._events if e.aggregate_type == agg_type and e.aggregate_id == agg_id][start_ver:]
for e in events:
version += 1
fn = globals().get(f"apply_{e.event_type}")
state = fn(state, e) if fn else state
if version % snapshot_every == 0: self._snapshots[key] = (version, state)
return state, version
def subscribe(self, event_type: str, fn: Callable[[DomainEvent], None]) -> None:
self._subs[event_type].append(fn)
def apply_OrderCreated(state, e): return {"order_id": e.aggregate_id, "items": {}, "status": "CREATED", **e.payload}
def apply_ItemAdded(state, e):
s = dict(state); pid = e.payload["product_id"]
s["items"] = {**s["items"], pid: s["items"].get(pid, 0) + e.payload["qty"]}
return s
def apply_OrderPlaced(state, e): return {**state, "status": "PLACED", "placed_at": e.occurred_at}
store = EventStore()
projection: dict[str, dict] = {}
store.subscribe("OrderPlaced", lambda e: projection.__setitem__(e.aggregate_id, {"status": "PLACED", **e.payload}))
CQRS:写读模型分离
写入侧走领域模型 + 事件,通过事件订阅异步构建读侧投影表(Projection)。写侧保证一致性,读侧针对查询做反范式和索引,各自最优,但引入最终一致性窗口。
| 维度 | 传统 CRUD | 纯 CQRS | CQRS + Event Sourcing |
|---|---|---|---|
| 存储 | 单表既读又写 | 写库+读库 | 事件流+N个投影表 |
| 查询能力 | 中等 | 极高 | 极高+时间旅行 |
| 审计追溯 | 额外加字段 | 难 | 天生免费 |
| 调试难度 | 低 | 中高 | 非常高 |
| 适合场景 | 大部分业务 | 读多写少查询复杂 | 金融/账本/溯源 |
最佳实践
先问三问再决定:业务是否需要不可变审计?查询是否真的复杂到读写冲突?团队是否有 DDD 经验?三个 Yes 才考虑。否则传统 CRUD + 审计表往往 ROI 更高。