Day 36 / 共 60 天 · 阶段 6 持久化与记忆
SqliteSaver:把内存字典换成真数据库
昨天 InMemorySaver 把档案存在内存字典里——进程一关就没了。今天看第一个真落盘的实现 SqliteSaver。你会发现"套路和 InMemory 一模一样,只是介质变了":三层字典变成两张 SQL 表、拆值存 blob 变成 INSERT 语句。同时它带来内存版没有的新问题:建表迁移、并发线程安全、事务。我们逐段读它 646 行里的核心,看清"一次 checkpoint 落盘到底跑了哪几条 SQL"。
📍 阶段 6 · 持久化与记忆(8 天)你在这里
D33 概念→
D34 接口→
D35 InMemory→
D36 Sqlite→
D37 Postgres→
D38 id体系→
D39 serde→
D40 Store
💡 用一个类比先兜住今天
InMemorySaver 是"把档案摊在你的办公桌上"——拿取飞快,但下班(进程退出)保洁一扫就没了。SqliteSaver 是把档案装进办公室角落的一个铁皮文件柜(一个
.sqlite 文件):柜子里有两个抽屉——checkpoints 抽屉放档案本体、writes 抽屉放半成品写入。第一次用要先"打柜子、贴好抽屉标签"(setup 建表)。多人同时开柜子会打架,所以配了一把锁(threading.Lock)。缺点:它是单机文件柜、不适合多人高并发——那要等明天的 Postgres。L01
存储从"三层字典"变成"两张表"
🤔 痛点:昨天 InMemory 用 storage/writes/blobs 三个字典。换到 Sqlite,数据结构怎么对应?
看类定义和它的两个核心属性(checkpoint/sqlite/__init__.py:45,82-95):
# sqlite/__init__.py:45
class SqliteSaver(BaseCheckpointSaver[str]):
conn: sqlite3.Connection # 数据库连接
is_setup: bool # 是否已建过表(懒建表用)
def __init__(self, conn, *, serde=None) -> None:
super().__init__(serde=serde)
self.jsonplus_serde = JsonPlusSerializer()
self.conn = conn
self.is_setup = False
self.lock = threading.Lock() # ← 线程锁,保证并发安全
BaseCheckpointSaver[str]和 InMemory 一样继承基类、泛型参数是 str——表示它用字符串版本号("数字.随机",Day 35 讲过)。conn / is_setup / lock三个家当:数据库连接、"建表了没"标志、一把线程锁。注意没有 blobs——Sqlite 版把通道值直接内联进 checkpoint 这一列(序列化后整包存),不像 InMemory 单独拆 blob。jsonplus_serde 另开一个除了继承来的 self.serde(存通道值用),还单独建一个 JsonPlus 序列化器——metadata 用 json.dumps 存成人类可读的 JSON(方便 SQL 里过滤查询)。🅰 设计取舍①:为什么 Sqlite 版把通道值内联进 checkpoint 列,不像 InMemory 那样拆 blob 去重?
InMemory 拆 blob 是为了"多份档共享没变的通道值"省内存(Day 35)。Sqlite 却选择把整个 checkpoint(含通道值)序列化成一个 BLOB 整包存进一列。原因:Sqlite 定位是轻量、单机、demo/小项目(docstring 明说 "lightweight, synchronous use cases")。整包存的实现最简单——一条 INSERT 搞定、读时一条 SELECT 拿回,不用像 Postgres 那样写复杂的 JOIN 去拼 blob。代价是长对话下存储会有冗余(没变的通道每份档都存一遍)。这是"实现简单性 vs 存储效率"的取舍:小项目不在乎那点冗余,要的是零心智负担。真要省空间和高并发,去用 Postgres(明天)。
L02
setup:建两张表 + 开 WAL 模式
第一次操作前要先建表。setup 用一段 SQL 脚本一次建好(checkpoint/sqlite/__init__.py:129-166):
-- sqlite/__init__.py:139 (self.conn.executescript 里)
PRAGMA journal_mode=WAL; -- ① 开启 WAL 日志模式
CREATE TABLE IF NOT EXISTS checkpoints ( -- ② 档案本体表
thread_id TEXT NOT NULL,
checkpoint_ns TEXT NOT NULL DEFAULT '',
checkpoint_id TEXT NOT NULL,
parent_checkpoint_id TEXT, -- 父档 id,串起家谱链
type TEXT,
checkpoint BLOB, -- 序列化后的档(含通道值)
metadata BLOB,
PRIMARY KEY (thread_id, checkpoint_ns, checkpoint_id) -- 三元组主键
);
CREATE TABLE IF NOT EXISTS writes ( -- ③ 半成品写入表
thread_id TEXT NOT NULL,
checkpoint_ns TEXT NOT NULL DEFAULT '',
checkpoint_id TEXT NOT NULL,
task_id TEXT NOT NULL,
idx INTEGER NOT NULL,
channel TEXT NOT NULL,
type TEXT,
value BLOB,
PRIMARY KEY (thread_id, checkpoint_ns, checkpoint_id, task_id, idx)
);
if self.is_setup: return方法开头就判断(:136)——只建一次。建过就直接返回,避免每次操作都跑建表。这叫"懒初始化"。PRAGMA journal_mode=WAL开 WAL(Write-Ahead Log 预写日志)模式:写操作先写日志、读操作不被写阻塞。让"一边写 checkpoint、一边读状态"能并发,大幅提升读写不打架的能力。checkpoints 表主键三元组(thread_id, checkpoint_ns, checkpoint_id) 联合主键——正是 Day 35 三层字典的三层键!字典嵌套在关系库里就表达成"联合主键"。writes 表主键五元组多了 task_id, idx——对应 Day 35 writes 字典的内层键 (task_id, idx)。主键唯一性正是"幂等去重"的数据库级保证(L06)。💡 CREATE TABLE IF NOT EXISTS 的深意加了
IF NOT EXISTS,建表语句可以安全地反复执行——表已存在就跳过、不报错。这让"懒建表"即使在多次调用/多进程下也不会因为"表已存在"而崩。这是一种朴素但极其实用的幂等 DDL:把"初始化"写成"运行多少次结果都一样",是数据库工程的基本功。L03
cursor:用一把锁串行化所有 DB 操作
所有 SQL 都通过一个统一的 cursor 上下文管理器执行(checkpoint/sqlite/__init__.py:168-189):
# sqlite/__init__.py:168
@contextmanager
def cursor(self, transaction: bool = True) -> Iterator[sqlite3.Cursor]:
with self.lock: # ① 先抢锁 —— 同一时刻只有一个线程能进
self.setup() # ② 确保表已建(懒建表就发生在这)
cur = self.conn.cursor()
try:
yield cur # ③ 把游标交给调用方执行 SQL
finally:
if transaction:
self.conn.commit() # ④ 默认提交事务(写操作要)
cur.close() # ⑤ 无论如何都关游标
with self.lock核心:每次 DB 操作都要先拿到那把 Lock。同一时刻只有一个线程能操作数据库——把并发访问"串行化",避免 sqlite3 连接被多线程同时使用出问题。self.setup()锁内调 setup——保证"建表"这个动作也是线程安全的(is_setup 标志的读写都在锁保护下)。transaction 参数写操作(put/put_writes)传 transaction=True,finally 里 commit() 落盘;纯读操作(get_tuple/list)传 False,不需要提交,省一次 commit 开销。finally 里 close无论 SQL 成功还是抛异常,游标都会被关闭——上下文管理器保证资源不泄漏。cursor 上下文管理器把"抢锁→建表→执行→提交→关游标"固化成统一流程
🅰 设计取舍②:为什么用一把粗粒度的全局锁,而不是更细的锁?
一把锁把所有数据库操作串成一队,并发度其实很低——但这正是 SqliteSaver 有意的选择。它连接时用
check_same_thread=False(:124)允许跨线程共享连接,代价就是必须自己保证串行访问,于是用锁兜底。为什么不做细粒度锁提升并发?因为 Sqlite 本身就不是为高并发写设计的(它是单文件、写时基本要独占)。花力气做细锁,底层 Sqlite 也扛不住并发写。所以"一把大锁 + 定位轻量场景"是自洽的:不假装能高并发,把简单和正确做到位。需要真并发?明天的 Postgres 用连接池 + pipeline。L04
put:一条 INSERT OR REPLACE 落盘
存一份档案,核心就一条 SQL(checkpoint/sqlite/__init__.py:387,418-443):
# sqlite/__init__.py:418
thread_id = config["configurable"]["thread_id"]
checkpoint_ns = config["configurable"]["checkpoint_ns"]
type_, serialized_checkpoint = self.serde.dumps_typed(checkpoint) # ① 整包序列化档
serialized_metadata = json.dumps(
get_checkpoint_metadata(config, metadata), ensure_ascii=False
).encode("utf-8", "ignore") # ② metadata 存 JSON
with self.cursor() as cur:
cur.execute(
"INSERT OR REPLACE INTO checkpoints "
"(thread_id, checkpoint_ns, checkpoint_id, parent_checkpoint_id, type, checkpoint, metadata) "
"VALUES (?, ?, ?, ?, ?, ?, ?)",
(
str(thread_id), checkpoint_ns, checkpoint["id"],
config["configurable"].get("checkpoint_id"), # ③ 父 id = 当前 config 的 id
type_, serialized_checkpoint, serialized_metadata,
),
)
return {"configurable": {"thread_id": thread_id, "checkpoint_ns": checkpoint_ns,
"checkpoint_id": checkpoint["id"]}} # ④ 返回新坐标
dumps_typed(checkpoint)把整个 checkpoint(含通道值,不拆 blob)序列化成 (类型标记, 字节)。类型标记单独存 type 列,方便读回时知道用什么反序列化(Day 39)。INSERT OR REPLACESqlite 特有语法:主键冲突就整行替换。这让 put 天然幂等——同一个 checkpoint_id 重复 put(比如恢复重跑),覆盖即可,不会插出两行。父 id = config 里的 checkpoint_id和 InMemory 一样的巧思:写新档时,config 里携带的 checkpoint_id 正是"上一份档",拿它当 parent_checkpoint_id,家谱链就串起来了。参数化 ? 占位全用 ? 占位 + 元组传参——防 SQL 注入的标准做法,值永远作为数据传入,不拼进 SQL 字符串。📝 走一遍
第 3 步存档:thread_id="chat-1"、checkpoint_ns=""、id="1efab...c3"、父 id="1efab...b2"。执行一条 INSERT OR REPLACE,checkpoints 表多一行;返回的 config 把 checkpoint_id 换成新的 "1efab...c3",供下一步当父。整个通道值(messages 等)全在那一列的 BLOB 里。
L05
get_tuple:一查档、二查半成品
取档比存档多一步——档和它的 pending_writes 在两张表,要查两次(checkpoint/sqlite/__init__.py:191,226-293):
# sqlite/__init__.py:226
checkpoint_ns = config["configurable"].get("checkpoint_ns", "")
with self.cursor(transaction=False) as cur: # 读操作,不开事务
if checkpoint_id := get_checkpoint_id(config): # 给了 id → 查指定档
cur.execute("SELECT ... FROM checkpoints WHERE thread_id=? AND checkpoint_ns=? AND checkpoint_id=?", ...)
else: # 没给 id → 查最新档
cur.execute("SELECT ... FROM checkpoints WHERE thread_id=? AND checkpoint_ns=? "
"ORDER BY checkpoint_id DESC LIMIT 1", ...) # ← 按 id 降序取第一条
if value := cur.fetchone():
(thread_id, checkpoint_id, parent_checkpoint_id, type, checkpoint, metadata) = value
# 第二次查询:捞这份档的 pending_writes
cur.execute("SELECT task_id, channel, type, value FROM writes "
"WHERE thread_id=? AND checkpoint_ns=? AND checkpoint_id=? ORDER BY task_id, idx", ...)
return CheckpointTuple(
config,
self.serde.loads_typed((type, checkpoint)), # ① 反序列化档(通道值就在里面)
json.loads(metadata) if metadata else {}, # ② metadata 反序列化
({...parent...} if parent_checkpoint_id else None), # ③ 父坐标
[(task_id, channel, self.serde.loads_typed((type, value))) # ④ 半成品
for task_id, channel, type, value in cur],
)
ORDER BY checkpoint_id DESC LIMIT 1不给 id 时取"最新档"——因为 checkpoint_id 单调递增(Day 35/38),降序排第一条就是最后写的。等价于 InMemory 里的 max(keys),只是交给 SQL 做。两次 execute先查 checkpoints 表拿档本体,再查 writes 表拿这份档的所有半成品写入。两张表分别查、在 Python 里组装成 CheckpointTuple。loads_typed((type, checkpoint))用存时记下的 type 标记 + 字节,反序列化回完整 checkpoint。通道值不用像 InMemory 那样 _load_blobs 逐个捞——因为它当初就是整包存的,一次性解开。ORDER BY task_id, idx查半成品时按 (task_id, idx) 排序——保证 pending_writes 的顺序稳定、可复现(和写入时的 idx 对齐)。Day 35 的三层字典 → Sqlite 的两张表;blobs 字典被"内联进 checkpoint 列"省略了
L06
put_writes:REPLACE 还是 IGNORE 二选一
存半成品写入时,它会根据"是不是特殊写"切换 SQL(checkpoint/sqlite/__init__.py:445,462-482):
# sqlite/__init__.py:462
query = (
"INSERT OR REPLACE INTO writes (...) VALUES (?, ?, ?, ?, ?, ?, ?, ?)"
if all(w[0] in WRITES_IDX_MAP for w in writes) # ① 全是特殊写 → 允许覆盖
else "INSERT OR IGNORE INTO writes (...) VALUES (?, ?, ?, ?, ?, ?, ?, ?)" # ② 含普通写 → 已存在则忽略
)
with self.cursor() as cur:
cur.executemany(
query,
[
(str(thread_id), str(checkpoint_ns), str(checkpoint_id),
task_id,
WRITES_IDX_MAP.get(channel, idx), # ③ 特殊写用负数 idx,普通写用序号
channel, *self.serde.dumps_typed(value))
for idx, (channel, value) in enumerate(writes)
],
)
WRITES_IDX_MAPDay 34/35 见过:{ERROR:-1, SCHEDULED:-2, INTERRUPT:-3, RESUME:-4}(checkpoint/base/__init__.py:795)。特殊写用负数 idx,普通通道写用它在列表里的正序号。全特殊 → INSERT OR REPLACE如果这批写全是特殊写(如 RESUME 恢复值),用 REPLACE 允许覆盖——恢复值该被最新的更新。含普通写 → INSERT OR IGNORE只要有普通写,用 IGNORE:主键已存在就静默跳过不覆盖。这就是 Day 35 InMemory "普通写幂等去重"在 SQL 层的等价实现——靠五元组主键唯一性 + IGNORE 兜底。executemany一批写入用 executemany 一次提交,比逐条 execute 高效。⚠ 边界:幂等靠"主键 + IGNORE",而不是先查再插
注意它没有"先 SELECT 看在不在、再决定插不插"——那样有竞态(查和插之间别的线程可能插了)。它直接
INSERT OR IGNORE:把幂等交给数据库主键约束,冲突了数据库自己忽略。这是数据库工程里的黄金法则:能用约束保证的正确性,绝不用应用层的"查了再改"去凑——后者几乎必有竞态。SqliteSaver 用一行 IGNORE 就把"重放不重复"做对了。L07
异步方法直接拒绝 + 小结
还记得 Day 35 InMemory 的异步方法"直接调同步"吗?SqliteSaver 更干脆——直接不支持异步(checkpoint/sqlite/__init__.py:585-624):
# sqlite/__init__.py:585
async def aget_tuple(self, config) -> CheckpointTuple | None:
raise NotImplementedError(_AIO_ERROR_MSG) # 抛异常,附带指路信息
# _AIO_ERROR_MSG(:34)大意:
# "SqliteSaver 不支持 async。请改用 AsyncSqliteSaver(需装 aiosqlite):
# from langgraph.checkpoint.sqlite.aio import AsyncSqliteSaver"
async def aput(self, config, checkpoint, metadata, new_versions):
raise NotImplementedError(_AIO_ERROR_MSG)
raise NotImplementedError异步方法直接抛错,附一段清晰的指路信息——告诉你去用 AsyncSqliteSaver。不假装支持、不静默降级。为什么不像 InMemory 那样"异步调同步"?因为 Sqlite 有真磁盘 IO。若在异步事件循环里同步调 sqlite(会阻塞 IO),会卡住整个 event loop、拖垮所有协程。所以它宁可报错,逼你用基于 aiosqlite 的真异步实现。get_next_version 同 InMemory版本号生成(:626-645)和 InMemory 完全一样:"32 位数字.16 位随机",复用同款"单调 + 抗碰撞"策略(Day 35 L06)。💡 本质:有真 IO 时,"同步实现塞进异步接口"是有害的Day 35 InMemory 无 IO,异步直接调同步没问题。Sqlite 有磁盘 IO,同步调用会阻塞 event loop——所以它选择"报错 + 指路"而非"凑合支持"。这是一条重要的工程判断:是否该提供某个接口,取决于能不能正确地提供,而不是能不能勉强跑通。明天的 Postgres 会展示"真异步该怎么认真实现"。
🧠 今日小结自测
- Sqlite 用几张表?分别对应 InMemory 的什么?(两张:checkpoints=档本体、writes=半成品;blobs 被内联进 checkpoint 列省略)
- setup 为什么用 IF NOT EXISTS?(幂等 DDL,可反复安全执行)
- WAL 模式解决什么?(读写不互相阻塞,提升并发读写能力)
- cursor 里的锁作用是什么?(把所有 DB 操作串行化,配合 check_same_thread=False 保证线程安全)
- put_writes 的幂等靠什么?(五元组主键 + INSERT OR IGNORE,让数据库约束保证不重复,避免应用层竞态)
- 为什么异步方法直接报错而不是调同步?(Sqlite 有真磁盘 IO,同步调用会阻塞 event loop)
✋ 10 分钟动手
# 1. 读 setup 建表 + cursor 锁
sed -n '129,189p' libs/checkpoint-sqlite/langgraph/checkpoint/sqlite/__init__.py
# 2. put / put_writes 两段对照读
sed -n '418,443p;462,482p' libs/checkpoint-sqlite/langgraph/checkpoint/sqlite/__init__.py
# 3. 亲手落盘一个对话,再用 sqlite3 命令行看两张表
python - <<'PY'
from langgraph.graph import StateGraph, START
from langgraph.checkpoint.sqlite import SqliteSaver
import sqlite3
g=StateGraph(dict); g.add_node("a",lambda s:{"n":s.get("n",0)+1}); g.add_edge(START,"a")
conn=sqlite3.connect("demo.sqlite", check_same_thread=False)
app=g.compile(checkpointer=SqliteSaver(conn))
app.invoke({"n":0}, {"configurable":{"thread_id":"t1"}})
print("表:", conn.execute("SELECT name FROM sqlite_master WHERE type='table'").fetchall())
print("档数:", conn.execute("SELECT count(*) FROM checkpoints").fetchone())
PY
🔮 明日预告 · Day 37 PostgresSaverSqlite 是单机文件柜、串行化、不支持异步。明天看生产级的
PostgresSaver:你会看到它真的把通道值拆进 checkpoint_blobs 表去重(回到 InMemory 的做法)、用 MIGRATIONS 版本化迁移管理表结构演进、用 JSONB + 一条带 JOIN 的 SELECT 在数据库里就把通道值拼好、还有 Pipeline 批量提交和连接池。它是"同样的套路,工业强度实现"。