第 13 期:换真实存储——SQLite、Postgres 与长任务
第 12 期把 checkpointer 和 store 落在两个 SQLite 文件里,这是本机跑通一切 最省事的选择。给别人用的服务不能停在这一步:真实部署里,进程会被重启、 会被换到另一台机器上、容器会被重新拉起——第 14 期就要把这个服务放到 Render 上。SQLite 是磁盘上的一个文件,这个事实在这一期第一次露出真正 的代价。这一期换成 Postgres,顺带补一条第 12 期没做的路:一个真的会跑 一阵子的任务,不该占着一条 HTTP 连接干等。
敲进去
第 13 期的代码在 code/ep13/,graph.py/tools.py/state.py 跟第 12 期
一字不改。
存储:接口没变,连接串变了
async def _open_storage():
if settings.POSTGRES_URL:
from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver
from langgraph.store.postgres.aio import AsyncPostgresStore
saver_cm = AsyncPostgresSaver.from_conn_string(settings.POSTGRES_URL)
store_cm = AsyncPostgresStore.from_conn_string(settings.POSTGRES_URL)
saver = await saver_cm.__aenter__()
store = await store_cm.__aenter__()
await saver.setup()
await store.setup()
return saver, store, saver_cm, store_cm
from langgraph.checkpoint.sqlite.aio import AsyncSqliteSaver
from langgraph.store.sqlite.aio import AsyncSqliteStore
saver_cm = AsyncSqliteSaver.from_conn_string(str(CHECKPOINT_DB))
store_cm = AsyncSqliteStore.from_conn_string(str(MEMORY_DB))
saver = await saver_cm.__aenter__()
store = await store_cm.__aenter__()
await store.setup()
return saver, store, saver_cm, store_cm
AsyncPostgresSaver/AsyncPostgresStore 跟 AsyncSqliteSaver/
AsyncSqliteStore 是同一套接口——from_conn_string()、.setup(),
连方法名都没变,只是连接串从一个文件路径换成一条 postgresql://...。
build_graph 拿到手的还是同一个 saver/store 对象,图那边一行代码
不知道、也不需要知道后面换了存储。没设 POSTGRES_URL 就照旧用 SQLite——
不逼着每个跟读的人先装 Postgres。
.setup() 建的表,真机跑一次能看到:
checkpoints、checkpoint_blobs、checkpoint_writes、checkpoint_migrations
store、store_migrations
长任务:提交和取结果,拆成两次请求
@app.post("/chat/submit", dependencies=[Depends(check_auth)])
async def submit(req: ChatRequest, request: Request) -> dict:
graph = request.app.state.graph
config = build_run_config(req.thread_id, req.user_id)
state = {"messages": [HumanMessage(req.message)]}
task = asyncio.create_task(graph.ainvoke(state, config=config))
request.app.state.tasks[req.thread_id] = task
return {"thread_id": req.thread_id, "status": "submitted"}
@app.get("/chat/{thread_id}/status", dependencies=[Depends(check_auth)])
async def status(thread_id: str, user_id: str, request: Request) -> dict:
graph = request.app.state.graph
config = {"configurable": {"thread_id": thread_id, "user_id": user_id}}
snapshot = await graph.aget_state(config)
if not snapshot.values:
raise HTTPException(status_code=404, detail="没有这个 thread_id")
interrupted = bool(snapshot.interrupts)
done = not snapshot.next and not interrupted
last = snapshot.values["messages"][-1]
return {
"thread_id": thread_id,
"status": "interrupted" if interrupted else ("done" if done else "running"),
"last_message": last.content if not getattr(last, "tool_calls", None) else None,
"interrupt_payload": snapshot.interrupts[0].value if interrupted else None,
}
/chat/submit 起一个后台 asyncio.Task 就立刻回,不等图跑完;
/status 现查 graph.aget_state()——这一步读的是 Postgres 里落盘的
最新一条记录,不是某个进程内存里的变量。request.app.state.tasks
只是拿住这个 Task 的引用,防止它被垃圾回收提前打断,本身不参与
状态查询。
跑起来
docker run -d --name pg -e POSTGRES_PASSWORD=devpass -e POSTGRES_DB=langgraph_ep13 -p 5433:5432 postgres:17
export POSTGRES_URL=postgresql://postgres:devpass@localhost:5433/langgraph_ep13
cd code
uv run uvicorn ep13.app:app --port 8000
curl -X POST http://localhost:8000/chat/submit \
-H "Authorization: Bearer sk-xxxx" -H "Content-Type: application/json" \
-d '{"user_id":"wang","thread_id":"t1","message":"帮我一起查一下 KL-778、KL-901、KL-315 这三个订单能不能改期"}'
curl "http://localhost:8000/chat/t1/status?user_id=wang" -H "Authorization: Bearer sk-xxxx"
你应该看到什么
长任务:提交立刻回,过一会儿状态变成 done
$ curl -X POST http://localhost:8000/chat/submit ... -d '{...三订单并行核对...}'
{"thread_id":"pg1","status":"submitted"}
$ curl "http://localhost:8000/chat/pg1/status?user_id=wang" ...
{"thread_id":"pg1","status":"running","last_message":"[并行核对结果]\nKL-778:...\nKL-901:...\nKL-315:...","interrupt_payload":null}
$ curl "http://localhost:8000/chat/pg1/status?user_id=wang" ...
{"thread_id":"pg1","status":"done","last_message":"三个订单的核对结果如下:\n\n- **KL-778**:可以改期,出行日前3天内可免费改期一次。\n- **KL-901**:出行日9月5日,需提前3天(即9月2日前)改期,今天已进入3天内,**不支持改期**。\n- **KL-315**:可以改期,任意时间免费改,但须在原出发时间前完成。\n\n需要我帮您操作哪个订单的改期吗?"}
第一次轮询正好撞在“并行核对已经跑完、模型还没把三份结果组织成最终
回答“这个中间点上——status 是 running,last_message 是第 10 期
那条 [并行核对结果] 的拼接消息(没有 tool_calls,符合“最后一条不带
工具调用的消息“这个判断,但还不是真正的最终答案)。第二次轮询才是
done,答案里还带着模型自己算出来的日期判断——KL-901 已经过了免费
改期窗口,这条推理和第 11 期真机撞见的一致。
interrupt 也能轮询到:
$ curl -X POST http://localhost:8000/chat/submit ... -d '{"message":"帮我取消订单 KL-901", ...}'
{"thread_id":"pg2","status":"submitted"}
$ curl "http://localhost:8000/chat/pg2/status?user_id=chen" ...
{"thread_id":"pg2","status":"interrupted","last_message":null,"interrupt_payload":{"action":"cancel_order","order_id":"KL-901","customer":"陈先生","product":"东京迪士尼一日票"}}
存储:容器重启,SQLite 忘光,Postgres 记得
先用第 12 期的 SQLite 版本,记一句话:
$ curl -N -X POST http://localhost:8002/chat ... -d '{"thread_id":"restart1","message":"我叫王小姐,我的订单是 KL-778,帮我记住"}'
data: {"type": "answer", "content": "已经帮您记住了,王小姐。..."}
把这个进程杀掉,挪走它的 SQLite 文件(模拟“这台容器没了,新容器的
本地磁盘是空的“),原地重新起一个一模一样的进程,问同一个 thread_id:
$ curl -N -X POST http://localhost:8002/chat ... -d '{"thread_id":"restart1","message":"我刚才说我叫什么、订单号多少?"}'
data: {"type": "answer", "content": "您好,这是我们这次对话的开头,您还没有告诉我您的姓名和订单号呢。..."}
完全不记得——这不是 bug,是 SQLite 文件天生的边界:它只活在写下它的 那台机器的本地磁盘上,换一块盘,之前的内容就是不存在。
换成这一期的 Postgres 版本,重复一模一样的实验——记一句话,杀掉进程,
POSTGRES_URL 不变地重新起一个新进程:
$ curl -N -X POST http://localhost:8003/chat ... -d '{"thread_id":"restart-pg","message":"我叫王小姐,我的订单是 KL-778,帮我记住"}'
data: {"type": "answer", "content": "好的,已记下:王小姐,常用订单号 KL-778。..."}
(杀掉整个进程,重新启动,进程 ID 变了,本地没有留下任何文件)
$ curl -N -X POST http://localhost:8003/chat ... -d '{"thread_id":"restart-pg","message":"我刚才说我叫什么、订单号多少?"}'
data: {"type": "answer", "content": "您刚才说自己是**王小姐**,常用订单号是 **KL-778**,已帮您记在档案里,用于身份核验。..."}
进程本身没有留下任何痕迹,记忆却完整地在——因为它从来就没有存在 “进程“这一层,存在的是 Postgres 里的两行数据。
--workers 2 也真机测过:两个真实的 uvicorn worker 进程共用同一个
POSTGRES_URL,四个并发提交的长任务全部轮询到 done,跟提交请求
落在哪个 worker、轮询请求又落在哪个 worker 完全无关。
发生了什么
SQLite 的边界不是“并发写会炸“,是“文件只活在一块盘上“。 这一期
真机测过 SQLite 在 --workers 多进程下的并发写入,没有炸——aiosqlite
的重试和 WAL 模式扛住了这本书这个量级的并发。真正测出来的边界是另一
件事:进程一旦真的重启(换了容器、换了机器、本地磁盘是新的),之前
写在这块盘上的一切都不存在了。生产环境里“进程被重启“是常态,不是
异常——容器编排系统调度、扩容、部署新版本,都会带来一次新的容器实例,
新实例的本地磁盘天然是空的。
Postgres 解决的不是“更快“或“更能扛并发“,是“存在的地方跟进程的 生死解耦“。 这一期没有做任何性能对比,两边的读写延迟在这个量级上 感觉不出差别。真正的差别是:Postgres 是一个独立运行的服务,进程连上 它、断开它,Postgres 自己的生死跟任何一个连接它的进程无关——这正是 “服务“和“给服务用的存储“必须分开部署的原因。
长任务这条支线,价值在“客户端可以走开“,不在“跑得更快“。 /chat/ submit 起的后台任务和 /chat 直接跑,花的时间一模一样——三个订单
该并行查的还是并行查,该等的网络往返一秒没少。变的是客户端这一侧:
不用为了等一个可能要跑十几秒的任务,占着一条 HTTP 连接、扛着中间可能
存在的负载均衡器超时设置。轮询和最初提交请求完全独立,甚至可以换一台
客户端设备去问“这个任务现在怎么样了“。
/status 能查到任意进程提交的任务,这件事本身就是“换真实存储“的
成果在起作用。 这一期没有专门为“跨进程共享任务状态“写任何代码——
snapshot = await graph.aget_state(config) 从第 4 期起就是这个签名,
这一期能查到别的 worker 提交的任务,纯粹是因为 aget_state 读的
Postgres 现在是所有 worker 共用的同一份,不是各自进程内存里的数据。
常见问题
CLI 版本(main.py)也换成 Postgres 了吗? 没有,main.py 还是
SQLite,这一期换的只是 app.py——命令行调试场景通常就是本机跑一次,
用不上跨进程共享,SQLite 更省事。
/chat/submit 起的 asyncio.Task,进程本身重启了会怎样? 这个
任务会跟着进程一起消失,status 会一直停在跑到一半的最后一个
checkpoint 上,不会自动变成 done,也不会报错——这是这一期故意没做
的部分,加分练习会补一种检测方式。
为什么 /status 不用 SSE,用轮询? 这条路径本来的目的就是“客户端
可以走开“,SSE 要求客户端保持连接,跟这个目标反着来。真要做“结果一
出来就主动推给客户端“,那是另一套机制(WebSocket、或者服务端主动
调用一个回调地址),这一期没做。
Postgres 连不上会怎样? AsyncPostgresSaver.from_conn_string()
连接失败会在 lifespan 里直接抛异常,FastAPI 应用整个起不来——这是
故意的:宁可看得见“启动失败“,也不要在存储层坏掉的情况下假装正常
服务。
加分练习
- 给
/chat/submit提交的任务加一个“最后活跃时间“记录,/status发现任务已经跑了很久还没完成、又查不到对应的asyncio.Task(说明提交它的那个进程已经不在了),返回一个第三种状态而不是 一直显示running。 - 把
POSTGRES_URL指向一个远端真实数据库(不是本机 docker), 量一下网络延迟对这一期几个接口的真实影响。 - 用
docker stop/docker start真的重启一次 Postgres 容器(不是 重启 FastAPI 进程),确认数据卷挂载正确的话,Postgres 自己重启 也不丢数据——这是“存储和服务分开部署“要成立的另一半前提。 - 读一下
checkpoint_blobs/checkpoint_writes这两张表里实际存的 是什么,对照第 4 期“checkpointer 存的到底是什么“那节内容,看 Postgres 版本和 SQLite 版本存的信息是不是完全对等。