用 SQLite 事件日志给工作流装上“黑匣子”

发布时间:2026/8/3 9:23:59
用 SQLite 事件日志给工作流装上“黑匣子” 从上一篇的纯状态机继续上一篇产出的miniflow.py定义了TripState、Event、Command与纯函数decide。本篇把它们当作既定输入不重写业务规则demo_102.py会导入Event和decide而数据库只负责可靠记录。上一篇的输出是内存中的事件序列和最终TripState本篇要把同样的序列变成可查询、可重放、可审计的持久化日志。为什么选 SQLite 而不是先上 Kafka持久化执行最关键的第一步不是吞吐量而是事务边界。单机 SQLite 提供真正的提交、唯一约束、崩溃恢复和清楚的文件形态足以把协议验证正确。事件日志采用只追加模型不更新旧事件不把最终状态覆盖到一行 JSON。追加日志保留因果过程后面才能重放、迁移和解释故障。等模型稳定后再把相同协议迁到 PostgreSQL 或分区日志远比一开始在分布式组件里调试语义容易。设计日志表和原子追加我们为每个工作流维护单调递增的seq。主键(workflow_id, seq)同时承担顺序约束和并发冲突检测。kind是事件类型payload使用稳定排序的 JSONcreated_at仅用于运维观察绝不能参与业务决策。WAL 模式允许读者不阻塞写者synchronousFULL更偏安全代价是每次提交可能触发刷盘。保存为event_store.py。这里有一个容易漏掉的细节连接级 PRAGMA 必须在每个新连接上设置不能以为初始化数据库时设置一次就永久有效。journal_modeWAL会持久化但foreign_keys、busy_timeout等很多选项不会。当前代码每次打开连接都走connect把行为集中起来。importjsonimportsqlite3fromcontextlibimportcontextmanagerfrompathlibimportPath SCHEMA CREATE TABLE IF NOT EXISTS events ( workflow_id TEXT NOT NULL, seq INTEGER NOT NULL, kind TEXT NOT NULL, payload TEXT NOT NULL, created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (workflow_id, seq) ); classEventStore:def__init__(self,path:str):self.pathPath(path)withself.connect()asdb:db.executescript(SCHEMA)contextmanagerdefconnect(self):dbsqlite3.connect(self.path,timeout5)db.row_factorysqlite3.Row db.execute(PRAGMA journal_modeWAL)db.execute(PRAGMA synchronousFULL)try:yielddb db.commit()exceptException:db.rollback()raisefinally:db.close()defappend(self,workflow_id:str,expected_seq:int,kind:str,payload:dict)-int:next_seqexpected_seq1bodyjson.dumps(payload,ensure_asciiFalse,sort_keysTrue)withself.connect()asdb:db.execute(INSERT INTO events(workflow_id, seq, kind, payload) VALUES (?, ?, ?, ?),(workflow_id,next_seq,kind,body),)returnnext_seqdefload(self,workflow_id:str)-list[dict]:withself.connect()asdb:rowsdb.execute(SELECT seq, kind, payload FROM events WHERE workflow_id? ORDER BY seq,(workflow_id,)).fetchall()return[{seq:row[seq],kind:row[kind],payload:json.loads(row[payload])}forrowinrows]运行输出模块定义成功无标准输出乐观并发不是可选装饰expected_seq看似多余却是避免“丢失推进”的关键。两个 worker 同时读到序号 2都想写序号 3复合主键只允许一个成功另一个收到IntegrityError后必须重新加载历史。不能用SELECT MAX(seq)1后无条件插入并期待数据库替你排队因为读和写之间存在竞态。更不能发生冲突后自动改写成序号 4第二个事件的决策基于旧状态直接顺延会把过期决定伪装成合法新事实。下面保存为demo_102.py。它复用上一篇miniflow.Event写入四个事件然后故意用过期版本再次追加以证明冲突会被拒绝。为便于重复运行示例使用临时目录真实服务则传入固定路径。importsqlite3importtempfilefrompathlibimportPathfromevent_storeimportEventStorefromminiflowimportEventwithtempfile.TemporaryDirectory()asdirectory:pathPath(directory)/flow.dbstoreEventStore(str(path))workflow_idtrip-001seq0history[Event(trip_requested,{city:成都}),Event(budget_locked,{token:B-7}),Event(flight_booked,{ref:F-8}),Event(hotel_booked,{ref:H-9}),]foreventinhistory:seqstore.append(workflow_id,seq,event.kind,event.data)print(appended,seq,event.kind)try:store.append(workflow_id,2,hotel_booked,{ref:stale})exceptsqlite3.IntegrityError:print(stale writer rejected)loadedstore.load(workflow_id)print(order:,[row[kind]forrowinloaded])assert[row[seq]forrowinloaded][1,2,3,4]运行输出appended 1 trip_requested appended 2 budget_locked appended 3 flight_booked appended 4 hotel_booked stale writer rejected order: [trip_requested, budget_locked, flight_booked, hotel_booked]事务边界到底保护了什么当前append一次只写一个事件所以原子性看起来理所当然。真正的工作流引擎通常还要同时写待执行命令也就是 outbox。若先提交事件再写命令中间崩溃会得到“状态已推进、活动却永远没有调度”的孤儿状态若先发送网络请求再提交事件会得到“副作用成功、历史却说没发生”的幽灵执行。第四篇会把事件和 outbox 放入同一事务。本篇先把日志协议单独钉牢是为了分清存储正确性与投递正确性。非平凡踩坑是把 SQLite 的“数据库已提交”误解成“任何硬件故障都绝不丢”。可靠性仍受文件系统、磁盘缓存、挂载参数和备份方式影响。WAL 数据可能暂时位于-wal文件只复制主.db文件会产生不一致备份。在线备份应使用 SQLite backup API或者在受控 checkpoint 后复制完整文件集合。容器中还要确认数据库位于持久卷而不是随 Pod 消失的临时层。另一个记忆点是序号表达的是工作流内因果顺序不是全球时间。不要用自增整数比较两个不同工作流谁先发生也不要拿created_at做重放排序。同一工作流依赖复合主键的严格序号跨工作流统计可以使用时间但它只是观测维度。把业务因果和墙上时钟混为一谈遇到时钟回拨、批量导入或跨地域复制时会付出代价。为恢复准备输入本篇产出了event_store.py的EventStore.append与EventStore.load并得到一个含四条顺序事件的flow.db。数据库现在能记住发生过什么却还不会从日志恢复TripState。下一篇将直接读取load返回的字典列表重新调用上一篇的decide证明删掉所有内存状态后仍能精确恢复并进一步解决“重放时不能重复产生外部命令”的问题。 觉得有用就点个赞 收藏方便回头查阅有疑问直接在评论区留言我看到都会回。 文章里的代码都能直接跑。想要可直接 clone 的完整工程 配套部署脚本 / 踩坑清单评论一声或发邮件到cj2664qq.com我免费发你。如果你正好在做类似系统、或有工程化难题想找人做也欢迎邮件聊一句——我按实际情况评估能落地的就接单或出方案。评论和邮件都能直接找到我不用跳别的平台。