CQRS + Event Sourcing 架构模式解析

Choyeon· 2026年9月14日· 3 分钟阅读· 244 阅读· 702 字· 2,423 字符· 更新于 2026年10月1日
CQRS + Event Sourcing 架构模式解析

当业务发展到一定规模,传统 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 更高。

本文作者

评论 (0)

暂无评论,来抢沙发吧。