尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Python轻量落地事件溯源与CQRS:事件存储到读模型实战

Python轻量落地事件溯源与CQRS:事件存储到读模型实战 做业务系统久了一定会撞上一个烦心事线上数据被改坏了或者某个状态莫名其妙跳变你却翻遍操作日志、数据库流水都拼不出它原本应该长什么样。我在维护交易类系统时不止一次经历这种“对账地狱”这时候才真正意识到“记录发生了什么”和“记录最终结果”是两码事。事件溯源Event Sourcing和CQRS恰恰是这两件事的正式化产物。但Python生态里聊这两个概念的人一直不多很多团队一听“事件溯源DDD落地”就被重量级标签吓退了。这篇文章想做的就是把这套思想放到Python熟悉的环境里用一张表、一个类、一个轻量消息通道把事件溯源和CQRS真正落下来。1. 为什么Python生态里事件溯源 CQRS 总被劝退1.1 事件溯源的核心思想只追加不修改要聊落地先得把概念掰扯清楚。事件溯源的核心是业务系统的最终真相不是数据库里的“当前状态”而是一串“已经发生过的事实”。订单从创建到支付到发货传统表里存的是 statuspending、statuspaid、statusshipped 这样的覆盖式更新最后只剩一个值事件溯源则把 OrderPlaced、OrderPaid、OrderShipped 这些事件原封不动地追加保存。当前状态只是事件流在某个时间点的一个投影随时可以从头重放出来。我用一个类比帮助你理解事件流就像一卷胶卷按快门的那一瞬被固定下来数据库表则像一张冲洗出来的照片只反映某一刻的样子。胶卷丢了照片还能再冲照片丢了胶卷还在就能重新洗。做业务审计时事件溯源的价值就是能精确回答“这单到底经历了什么”而不是靠应用日志去猜。1.2 Python项目的现实痛点数据没有“因果链条”Python后端的主流场景里Django、Flask、FastAPI 搭配一套ORM就能把CRUD跑起来。问题出在业务开始变复杂之后订单、支付、库存、物流各自写在不同的Service里互相调来调去每个模块都直接操作自己的数据表。一旦出现跨模块的状态不一致你根本说不清是谁先改了谁。我在一个电商中台项目里就遇到过这种事用户下单后订单表显示已支付但支付流水表里却没有对应记录。因为支付回调里先更新了订单表再写流水表中间进程重启了两条写操作没落在同一个事务里。这类问题常规代码很难防住因为“先做什么后做什么”只存在于代码逻辑里没有沉淀成可追溯、可回放的数据。CQRS在这里的价值不是“高大上”而是把“写”和“读”从一张表里拆开命令侧负责校验和业务规则产生明确的事件查询侧只负责把事件投影成好查询的读模型。读写各自优化问题边界清晰得多。1.3 “轻量化”到底是什么意思我在网上搜“Python 事件溯源 CQRS”跳出来的基本都是Java那套重型框架要么是完整的事件存储服务要么是带Actor模型的持久化引擎。很多Python团队一看到引入成本就打了退堂鼓。轻量化的思路是不造轮子不引重框架用Python开发者本来就熟悉的几个库——Pydantic 做事件建模、SQLAlchemy 做事件存储、Redis Streams 做事件分发——组合出一套最低限度可用的事件溯源 CQRS。这套方案跑在单体或接近单体的服务里完全够用后期要拆分布式也能平滑升级。2. 事件模型设计如何用Pydantic把“事实”钉死2.1 先分清楚命令和事件落地事件溯源第一步不是建表而是区分两个概念命令Command是想做的事事件Event是已经发生的事。命令可能失败事件永远成功。命名上我会刻意给事件用过去式OrderPlaced、OrderPaid、OrderShipped。别小看这个习惯团队里只要有人混着用代码很快会乱套。命令是入口比如 CreateOrderCommand 带商品列表、用户ID、金额OrderPlaced 是命令执行完之后沉淀下来的事实。命令允许校验不过被拒绝事件一旦入库就代表业务已经跨过了某个不可变节点。2.2 用Pydantic定义不可变事件在Python里定义事件我会优先选Pydantic而不是普通dataclass。原因有三个字段校验开箱即用、序列化反序列化干净、后面配合版本兼容很顺手。更关键的是 Pydantic 的frozenTrue可以强制不可变从模型层面堵住“改历史事件”的手。事件需要携带什么字段我的原则是尽量小而完整。事件本身就要能回答“发生了什么”不要让消费者拿到事件后还得去查别的表才能理解业务。from pydantic import BaseModel, Field from datetime import datetime from uuid import uuid4 class Event(BaseModel): model_config {frozen: True} event_id: str Field(default_factorylambda: str(uuid4())) occurred_at: datetime Field(default_factorydatetime.utcnow) property def event_type(self) - str: return self.__class__.__name__ class OrderPlaced(Event): order_id: str customer_id: str items: list[dict] total_amount: float class OrderPaid(Event): order_id: str paid_amount: float payment_method: str class OrderShipped(Event): order_id: str tracking_number: str这里有个实操细节event_type直接用类名。好处是消费者代码里isinstance判断就行后续序列化时也方便在 payload 里直接存类名重放事件时再用event_type反射出类。不要把事件设计成一个大字典字段塞满业务数据后期所有消费方都在解析这个字典改一次字段名就要动一堆代码。2.3 事件版本与兼容策略业务上线之后事件结构几乎必然要加字段。轻量化的做法不是搞一套复杂的事件版本管理协议而是在 payload 里加一个version字段标注事件结构的版本号。新字段加进来时老事件不改新事件带上新版本号。消费者端用version判断遇到老版本事件就走老逻辑遇到新版本就走新逻辑。这种双轨兼容的成本很低但千万不要直接去改已经落库的历史事件。后面我会专门说为什么历史事件不能改。3. 聚合根与事件存储一张表 一个类跑通整套逻辑3.1 聚合根是状态机的剧本领域驱动设计里的聚合根Aggregate听起来很玄落到代码里就是“一个负责保护业务不变量、把内部变化以事件形式流露出来的对象”。所有状态的改变都必须经过聚合根不允许业务逻辑在聚合根外部随意改字段。我拿下单场景举例。OrderAggregate 持有订单状态place、pay、ship方法内部做校验通过校验后就产生对应事件。注意聚合根内部不直接操作数据库它只管产生事件和用事件重建状态。class OrderAggregate: def __init__(self, order_id: str): self.order_id order_id self.version 0 self.status pending self.items [] self.paid_amount 0.0 def apply(self, event: Event): if isinstance(event, OrderPlaced): self.status placed self.items event.items elif isinstance(event, OrderPaid): self.status paid self.paid_amount event.paid_amount elif isinstance(event, OrderShipped): self.status shipped def place(self, items: list[dict], total_amount: float): if self.status ! pending: raise ValueError(订单当前状态不能下单) return OrderPlaced( order_idself.order_id, itemsitems, total_amounttotal_amount, ) def pay(self, amount: float, method: str): if self.status ! placed: raise ValueError(订单不在待支付状态) return OrderPaid( order_idself.order_id, paid_amountamount, payment_methodmethod, )看到apply(event)这个方法没有它是整个事件溯源的核心接口。聚合根有两个“出生方式”新建时从空白状态开始place产生第一个事件从存储恢复时则一条条apply历史事件把状态一步步重建出来。这两条路径走到最后聚合根的状态一定一致因为重建逻辑就是业务逻辑本身。3.2 事件存储表结构设计事件要存到哪里用一张专门的event_store表就够了不需要引入专门的EventStore服务。字段类型说明idbigint PK自增主键用于物理排序stream_idvarchar聚合根ID比如 order_idversionint聚合根版本号从1开始event_typevarchar事件类名用于反序列化payloadjsonb事件所有字段的JSON序列化created_attimestamp事件发生时间stream_id和version要建联合唯一约束这是后续乐观锁并发控制的基础。为什么用jsonb存 payload因为事件结构会演进关系型字段一概死板JSONB既能查询又能避免反复迁移表结构。MySQL用户用 JSON 类型也差不多。3.3 保存与重建的完整链路事件舱的核心操作只有两个append追加事件、load重建聚合根。我把两段核心逻辑写在这里。# 追加事件时必须带 expected_version用乐观锁防并发覆盖 def append(session, stream_id, events, expected_version): current_version session.scalar( select(func.max(EventRecord.version)) .where(EventRecord.stream_id stream_id) ) or 0 if current_version ! expected_version: raise ConcurrentModificationError( fstream {stream_id} 版本冲突期望 {expected_version}实际 {current_version} ) for idx, event in enumerate(events, start1): session.add(EventRecord( stream_idstream_id, versionexpected_version idx, event_typeevent.event_type, payloadevent.model_dump_json(), ))# 重建聚合根时按 version 升序读取事件 def load(session, order_id: str) - OrderAggregate: records session.scalars( select(EventRecord) .where(EventRecord.stream_id order_id) .order_by(EventRecord.version) ).all() aggregate OrderAggregate(order_id) for record in records: event_class EVENT_REGISTRY[record.event_type] event event_class.model_validate_json(record.payload) aggregate.apply(event) aggregate.version record.version return aggregate这里有两个容易踩的坑。第一个是load时必须按version升序不能按id升序。虽然多数情况下两者一致但一旦做过数据迁移或批量导入id顺序就可能和version顺序脱节。第二个是 EVENT_REGISTRY 的注册表必须维护好反序列化事件时要能从event_type找到对应的Pydantic类否则老数据重放时直接报错。4. CQRS落地写命令与读模型真正分家4.1 命令侧校验集中到业务层CQRS 的“写”那一侧我主张的命令处理流程是接收命令 - 加载聚合根 - 聚合根方法做校验 - 产生事件 - 追加事件 - 发布事件。业务校验全部放在聚合根方法里而不是放在ORM模型或ViewSet里。这么说吧ORM模型是给数据库打工的聚合根是给业务打工的。两者职责一混就会出现“数据库字段校验和业务规则校验纠缠不清”的烂摊子。订单能否被支付、能否被取消这些规则只属于聚合根。把规则放到外部service就会面临多个地方各写一份今天这里改明天那里漏。def handle_create_order(cmd: CreateOrderCommand): aggregate OrderAggregate(cmd.order_id) event aggregate.place(cmd.items, cmd.total_amount) event_store.append(cmd.order_id, [event], expected_version0) publish_event(event)上面这段代码是整个命令侧的缩影只有四行。因为复杂业务大多藏在聚合根内部命令处理函数只负责“加载、执行、保存、发布”四个动作。只要代码库中大多数命令处理器都是这个模式写部分的整体逻辑就会非常可控。4.2 查询侧投影把事件流变成读模型CQRS 的“读”那一侧核心产物是“投影”Projection。投影做的事情很简单监听事件更新一张专供查询用的表。这张表和业务表不同它完全是按照查询需求设计的可以宽表、可以冗余、可以不满足三范式。还是拿订单举例。查询侧有一张order_list_view表字段包括订单号、客户名、商品摘要、状态、金额、下单时间、付款时间。这张表的更新逻辑订阅 OrderPlaced、OrderPaid、OrderShipped 三个事件事件来了就对相应记录做 insert 或 update。class OrderListViewProjection: def handle(self, event: Event): if isinstance(event, OrderPlaced): upsert_order( order_idevent.order_id, customer_idevent.customer_id, items_summarysummarize_items(event.items), total_amountevent.total_amount, statusplaced, ) elif isinstance(event, OrderPaid): update_order( order_idevent.order_id, statuspaid, paid_timeevent.occurred_at, ) elif isinstance(event, OrderShipped): update_order( order_idevent.order_id, statusshipped, tracking_numberevent.tracking_number, )你有没有发现投影逻辑和聚合根的apply逻辑很像区别只在于apply重建的是业务聚合根投影生成的是查询模型。这块逻辑是CQRS读写分离真正发生的地方——写侧继续用事件溯源读侧不再触碰事件存储而是消费事件生成高性能查询数据。4.3 同步投影与异步投影怎么选投影可以同步也可以异步。这是落地CQRS时最常见的一个选择题。我自己的经验是单体应用、业务量不大、QPS 三位数以内优先同步投影。做法是在append事件成功后同一个事务里更新投影表。这样读写两边看到的最终状态是一致的没有延迟窗口运维难度低。当命令频率高起来或者投影计算量变大比如要根据事件更新一个全量汇总报表同步投影会拖慢写接口这时候再切异步——事件发布到消息通道后台消费者处理投影。异步的代价是查询数据有一段延迟对很多报表场景来说完全可接受。维度同步投影异步投影一致性强一致写完立刻可查最终一致有延迟窗口写接口延迟受投影耗时影响投影耗时不影响写接口复杂度低无需消息通道需要消息通道和消费端幂等适用场景小规模单体、后台管理高频写入、跨服务消费5. 事件发布的保命设计同事务写Outbox再分发5.1 为什么保存完事件不能直接发消息事件存好了下一步肯定是把它通知给其他服务或投影消费者。很多人的第一反应是append成功之后立刻往Redis Streams或RabbitMQ发一条。这个做法在崩溃场景下会出大事。假设数据库写入成功消息发送时进程崩溃或网络超时“消息没发出去”叠加“数据库已写入”下游永远不知道这笔订单发生了。反过来如果先发消息后写数据库消息发出去了数据库写入失败下游却收到了一条不存在的事件。这两种情况都会破坏系统一致性而且极难排查。那些做过分布式系统的朋友应该已经猜到了解法就是事务性发件箱Transactional Outbox把“待发布的事件”和事件本身放进同一个数据库事务再用一个后台发布器扫描发件箱把事件发到消息通道发成功后标记完成。5.2 Outbox实现细节落地时不复杂。在append事件的同一段代码里同时往outbox表插入一条记录。两者在同一个数据库事务里要么一起成功要么一起失败。def append_and_queue(session, stream_id, events, expected_version): append(session, stream_id, events, expected_version) for event in events: session.add(OutboxRecord( event_idevent.event_id, stream_idstream_id, event_typeevent.event_type, payloadevent.model_dump_json(), statuspending, created_atevent.occurred_at, ))后台发布器的工作逻辑也很直接轮询outbox表中statuspending的记录按创建时间升序取出逐条发到消息通道成功后把status改为published或者直接删除记录。def publish_outbox_messages(): pending_messages get_pending_outbox_records(limit100) for msg in pending_messages: redis.xadd(event_stream, { event_id: msg.event_id, stream_id: msg.stream_id, event_type: msg.event_type, payload: msg.payload, }) mark_published(msg.id)注意这个发布器要和业务进程分隔开最好单独部署或者至少放到单独的worker进程中。原因很简单业务进程如果满载甚至崩溃发布器不能被拖下水同时发布器的日志、监控、重试策略都和业务接口不一样没必要耦合在一起。5.3 消息通道选型Redis Streams 最省心Python轻量化方案里消息通道我首推 Redis Streams其次 RabbitMQ最后才是 Kafka。很多人一听到“分布式事件”就想到Kafka但轻量化落地不需要重型基础设施。Redis Streams支持消费者组、消息持久化、失败重试而且Redis在绝大多数Python后端项目里本来就有部署成本为零。RabbitMQ的定位是任务分发使用前还要安装服务端。Kafka更重适合跨团队的大规模事件总线场景轻量项目引入明显过度。方案部署成本消费者组消息持久化适合场景Redis Streams低已有实例支持支持单体内部事件分发RabbitMQ中支持支持异步任务、多服务通道Kafka高支持强海量事件、跨团队总线5.4 消费端幂等至少一次语义下的兜底消息系统最常见的投递语义是“至少一次”也就是说同一条事件很可能被投递两次以上。这本身不是Bug但如果投影逻辑不能处理重复事件就会出现数据重复、统计翻倍。幂等方案不外乎两种。第一种是在投影表里记录last_processed_event_id重复事件到达时直接跳过。第二种更稳妥给事件的消息ID建唯一约束消费时先尝试插入冲突说明已处理过。我习惯两者结合Redis Streams 消费者组天然按事件ID追踪消费进度投影表的 upsert 操作按order_id event_id判断是否已应用过。这层保护一旦做好即使发布器重试、消费者重启也不会导致读模型错乱。6. 避坑实测版本冲突、事件重放与历史事件不能改6.1 并发写入的版本号冲突事件溯源架构里同一聚合根的并发写入是一个绕不开的坑。两个请求同时支付同一个订单都从load读到了版本号5都通过校验并各自生成了OrderPaid如果没有任何保护事件存储里就会出现两个版本冲突的事件。解决办法就是我前面讲的乐观锁。append方法里的expected_version判断不能只查出来在Python里做比较那样会留出竞态窗口。正确的做法是依赖数据库层面的事务隔离机制或者用条件更新语句让数据库保证原子性。我在生产环境里用的是带WHERE stream_id? AND max(version)?的原子更新影响行数为0就抛并发异常。异常处理到命令层后标准的策略是重新加载聚合根、重新执行业务校验、再尝试提交。对绝大多数业务来说重试一次就能解决因为真正的并发写冲突发生频率很低。6.2 重放事件的顺序与幂等性实现事件溯源几乎必然碰到“重放”需求投影表坏了要重建、新加了一个投影要灌历史数据、调试时要把聚合根恢复到某个历史状态。重放的第一准则是严格按version升序。聚合根内部apply方法天然依赖顺序如果你把版本反转的事件喂进去得到的“当前状态”完全是错的。我在一次投影迁移中吃过这个亏因为一张历史表里id和version不一致直接按id重放结果所有历史订单的status全部错乱最后只能删表从零重建。重放的第二准则是投影操作要幂等。重建读模型本质上是“把事件流重新播放一遍”每一条事件都可能被处理多次。投影的更新操作最好设计成“按order_id upsert”不要用“累加”这种天然不幂等的操作。比如统计用户累计下单金额重放事件时如果直接累加重复播放一次金额就翻倍了。6.3 “历史事件永远不要改”不是口号做事件溯源最需要刻在脑子里的纪律已落库的历史事件永远不要直接修改。哪怕你发现事件里的字段定义错了、数值错了也绝不能 UPDATE。原因很直接事件流是事实记录其他服务可能已经消费了这些事件投影表已经基于它们建好。你改一条历史事件所有消费方的状态都会偏离事实而重放逻辑又会让这种偏离变得不可控。正确的姿势是发补偿事件。比如订单金额算错了不要改 OrderPlaced 的金额字段而是追加一个 OrderAmountCorrected 事件携带修正后的金额。聚合根apply这个事件时更新状态投影也相应更新读模型。补偿事件让事实链完整保留所有消费者都能看到“原先错了后来纠正了”的整个过程这个价值在审计场景里无可替代。7. 从单体到分布式轻量化方案的演进路径很多人担心这套轻量方案将来撑不住了怎么办。我的看法是不要一开始就为不存在的分布式买单。事件溯源和CQRS本身的扩展路径是清晰的单体阶段同步投影Redis Streams可以支撑到每天几百万事件再往上走把 Redis Streams 换成 Kafka投影服务独立成单独模块其他服务通过消息订阅事件流。关键区别只在于事件存储和消息通道的规模而不是业务逻辑设计。聚合根、事件模型、投影逻辑这些核心代码从单体迁移到分布式时几乎可以原样保留。这也是事件溯源架构比较舒服的地方——业务层和数据流的设计天然支持水平扩展。我在实际项目中落地事件溯源也不是一蹴而就的。第一版只加了一张event_store表和几十行聚合根代码先解决了“审计无据”这个最痛的问题第二版才引入 CQRS 投影把报表查询从业务表上剥离出去第三版才加了 Outbox 发布器把事件同步给下游服务。每走一步都验证了价值每走一步风险都可控。如果你所在的团队正被数据对账、审计缺失、跨模块状态混乱这些问题困扰我建议从最小最轻的那个点切入选一个业务边界清晰的聚合根先让事件流跑起来。你会在第一次通过重放事件流还原出线上问题的那一刻彻底理解这套模式为什么值得落地。
返回列表