Keyboard shortcuts

Press or to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

第 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/AsyncPostgresStoreAsyncSqliteSaver/ 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需要我帮您操作哪个订单的改期吗?"}

第一次轮询正好撞在“并行核对已经跑完、模型还没把三份结果组织成最终 回答“这个中间点上——statusrunninglast_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 应用整个起不来——这是 故意的:宁可看得见“启动失败“,也不要在存储层坏掉的情况下假装正常 服务。

加分练习

  1. /chat/submit 提交的任务加一个“最后活跃时间“记录,/status 发现任务已经跑了很久还没完成、又查不到对应的 asyncio.Task (说明提交它的那个进程已经不在了),返回一个第三种状态而不是 一直显示 running
  2. POSTGRES_URL 指向一个远端真实数据库(不是本机 docker), 量一下网络延迟对这一期几个接口的真实影响。
  3. docker stop/docker start 真的重启一次 Postgres 容器(不是 重启 FastAPI 进程),确认数据卷挂载正确的话,Postgres 自己重启 也不丢数据——这是“存储和服务分开部署“要成立的另一半前提。
  4. 读一下 checkpoint_blobs/checkpoint_writes 这两张表里实际存的 是什么,对照第 4 期“checkpointer 存的到底是什么“那节内容,看 Postgres 版本和 SQLite 版本存的信息是不是完全对等。