先交代一个背景:上个月我接手了一个客服工单流转的 Agent 项目,要求从“能跑”升级成“能断了再跑”。之前那套代码是我早期用 while 循环手搓的,LLM 调用、工具调用、状态拼接全塞在一起。表面上看没毛病,但真正压测和灰度的时候,问题全冒出来了:进程一重启,所有进行中的对话直接归零;有个退款审批节点需要人工确认,硬是塞不进循环里;两个用户同时触发流式任务,状态直接互相污染。折腾了两周之后,我决定把整套执行逻辑迁到 LangGraph 上,用 PostgreSQL 作为 Checkpoint 后端做状态持久化,再用 AG-UI 协议把运行时事件暴露给前端。这篇文章就把这次迁移里我踩过的坑、想明白的原理、可以直接抄的配置全部记录下来。
这次的改造主线其实特别清晰:把“图执行器”变成“可断点续跑的运行时”。LangGraph 负责有向图的调度和中断,PostgreSQL Checkpoint 负责把每一步的状态落盘,AG-UI 负责把运行时的内部事件标准化地吐给外部系统。三者各管一段,组合起来就是一个完整的可恢复 Runtime。下面我从头到尾拆一遍,代码都是可以直接跑的那种。
1. 为什么手写 Loop 撑不住生产场景
1.1 手写循环的典型形态与隐性问题
先看我最早那版核心代码,相信很多人写过一模一样的:
def run_agent(user_message: str): messages = [{"role": "user", "content": user_message}] for step in range(MAX_STEPS): response = llm.invoke(messages) if not response.tool_calls: return response.content messages.append(response) for tool_call in response.tool_calls: tool_result = execute_tool(tool_call["name"], tool_call["args"]) messages.append({"role": "tool", "tool_call_id": tool_call["id"], "content": tool_result})这个版本在 Demo 阶段完全够用,逻辑直接、调试简单。但放到生产环境,它有几个非常致命的结构性缺陷。
第一,状态全在内存变量里,进程死了就没了。第二,循环的中断只能靠抛异常,没有办法“停住->等待外部输入->再继续”。第三,多用户并发时 messages 列表是共用的,必须加锁或者每个用户开一个进程,资源开销直接爆炸。第四,没有任何可观测性,你不知道 Agent 内部在第几步、卡在哪个工具调用上、历史状态是什么样。
这些问题的本质是:手写循环把“控制流”和“状态存储”耦合死了。控制流是顺序执行的,状态是易失的,两者都无法支撑复杂的生产场景。
1.2 从 Loop 到 Runtime 的思维转变
后来我想明白一件事:Agent 的核心不应该是“循环调用 LLM”,而应该是一个“可暂停、可恢复、可观测的执行运行时”。这个运行时需要三个基本能力:
- 状态持久化:每一步执行结果都能写到外部存储,进程崩溃后可以恢复到最近状态。
- 暂停与恢复:执行到人工审批、外部确认这类节点时,可以主动挂起,等外部信号到达后再继续。
- 外部可观测:前端或调用方能看到 Agent 当前在做什么、需要什么,而不是对着一个黑盒子空等。
LangGraph 本质上就是围绕这套思路设计的,它把执行流程抽象成图,把状态管理交给 Checkpoint 机制,把中断恢复做成第一公民特性。与其自己从零写一套 state machine,不如直接用经过验证的框架。
1.3 顺带说清楚 LangChain 和 LangGraph 的分工
很多朋友会混淆 LangChain 和 LangGraph,其实我刚开始也花了一段时间才理顺。LangChain 是一个组件库,负责封装模型调用、Prompt 模板、工具抽象、RAG 检索这些基础能力;LangGraph 是一个编排引擎,负责把各个组件用图的方式组织起来,并且处理状态流转、持久化、多 Agent 协作这些上层问题。
实际项目中,两者不是二选一的关系。你在 LangGraph 的节点内部,依然会用 LangChain 的ChatOpenAI或者工具装饰器;LangGraph 负责图和状态,LangChain 负责模型和工具的具体实现。理解了这层分工,后面配置起来就不会绕。
2. LangGraph 核心机制拆解:状态、节点、Checkpointer
2.1 StateGraph 的运转方式
LangGraph 构建应用的基本单位是StateGraph。它的核心抽象是:一个全局状态对象(State),一堆节点(Node),以及节点之间的条件边(Edge)。每次图被调用时,状态对象依次流过各个节点,节点从状态里读数据、做处理、返回更新,然后图执行器把返回值合并回全局状态。
这里有个关键设计:状态更新是按“键”合并的,不是整体覆盖。每个节点返回的是一个 dict,只会覆盖返回 dict 里包含的那些字段,其他字段保持不变。这设计让多个节点可以各管各的字段,互不干扰。对应到代码里:
from typing import TypedDict, Annotated class AgentState(TypedDict): messages: Annotated[list, operator.add] current_order: str approval_status: str注意messages用了operator.add做累加器,表示新的消息会追加到列表里,而不是整个替换。这是 LangGraph 里最常用也最容易忽视的细节,特别是从手写循环迁移过来的人,一开始经常在这里翻车。
2.2 忽略 Checkpointer 等于没用 LangGraph
LangGraph 的图本身只是定义了计算流程,真正让它变成 Runtime 的是checkpointer。我是这么理解的:图是“怎么做”,Checkpoint 是“做到哪了”。每一个节点执行完之后,LangGraph 会把当前状态快照、图执行位置、待恢复的数据一起写到 checkpointer 里。之后如果进程重启、容器被重新调度,只要还有同一个thread_id,图就能从最后一次 Checkpoint 继续跑。
这个机制特别像数据库的 WAL 日志:每一步都先记日志,再继续执行。LangGraph 的 Checkpoint 也是这个逻辑,它记录的不仅有状态值,还有图内部的执行栈。这也是它能做到“中断恢复”的根本原因,因为恢复的不只是数据,还有执行到了哪个节点的控制权。
如果没配 checkpointer,LangGraph 就是一个带条件的循环而已,完全体现不出它的优势。所以我的建议很直接:只要是上生产的 LangGraph 应用,第一件事就是配好持久化 Checkpoint。
2.3 三款主流 Checkpointer 的选型对比
LangGraph 内置了几种 Checkpoint 实现,各有各的适用场景。我把它们的优缺点整理成了一张表:
| 存储方案 | 适用场景 | 优点 | 缺点 |
|---|---|---|---|
内存MemorySaver | 开发调试、单进程演示 | 零配置、速度快 | 进程重启即丢失,多副本无法共享 |
SQLiteSqliteSaver | 单机部署、小规模并发 | 文件即存储,轻量可靠 | 多进程并发写容易锁库,高并发吃力 |
PostgreSQLPostgresSaver | 生产环境、多实例部署 | 事务支持完善,连接池成熟,天然适合多进程 | 需要额外维护数据库,配置稍重 |
我的选择很明确:本地联调用 SQLite,测试环境和生产环境一律用 PostgreSQL。原因很简单,Agent 运行时的状态就是业务数据的一部分,不能丢。PostgreSQL 的 ACID 事务、行级锁、连接池生态都比 SQLite 更适合做多实例共享的存储层。我们在压测环境里用 SQLite 跑 8 路并发就会出现database is locked,换到 PostgreSQL 之后,50 路并发都很稳定。
2.4 PostgreSQL 里到底存了些什么
用 PostgresSaver 初始化之后,数据库里会自动生成三张核心表,理解它们的用途对排障很重要:
checkpoints:Checkpoint 的主记录表,每条记录对应一次状态快照,包含thread_id、checkpoint_id、parent_checkpoint_id、图的执行类型等元数据。它是整个状态历史的目录索引。checkpoint_blobs:真正存放序列化状态数据的地方。LangGraph 会把 Python 对象序列化成 JSON 或字节流,写进这张表的blob字段。这是最占空间的表。checkpoint_writes:写入缓冲表。当一个节点内部有多个并行分支时,LangGraph 会把每个分支的中间结果临时写到这里,等所有分支完成后再合并进正式的 Checkpoint。这也是 LangGraph 支持并行分支的底层基础。
这三张表配合thread_id字段,就形成了“一个线程一部历史”的存储模型。你可以把每个thread_id当成一个独立的对话会话或者任务实例,它的完整生命周期都会记录在这三张表里。
2.5 初始化 PostgreSQL Checkpointer 的正确姿势
我见过不少人在这里栽跟头,包括我自己第一次用的时候也踩了。直接看代码:
from langgraph.checkpoint.postgres import PostgresSaver # 方式一:直接用连接串(适合脚本和短任务) checkpointer = PostgresSaver.from_conn_string( "postgresql://agent_user:agent_pass@localhost:5432/agent_db" ) # 方式二:生产推荐,用连接池避免频繁建连 from psycopg_pool import ConnectionPool pool = ConnectionPool( conninfo="postgresql://agent_user:agent_pass@localhost:5432/agent_db", max_size=20, ) checkpointer = PostgresSaver(pool) # 首次使用需要初始化表结构 checkpointer.setup()这里有两个必须注意的细节。第一,setup()必须在建表之前调用一次,它是幂等的,重复调用没问题;如果跳过这一步直接compile(),运行时会报找不到checkpoints表。第二,from_conn_string每次都新建数据库连接,在一个常驻进程里用它来跑生产流量,连接数会越积越多,所以我后来换成了psycopg_pool的连接池方式,连接复用之后性能和稳定性都上了一个台阶。
3. 中断与恢复:让 Agent 在关键节点“停下来等人”
3.1 为什么要用 interrupt 而不是手动存数据库
很多从手写循环迁移过来的开发者,会本能地想:我直接在每个需要人工介入的节点,把状态写到数据库表里,然后轮询数据库看用户有没有确认,不也能实现暂停和恢复吗?
理论上可以,但这种做法有几个严重问题:你要自己定义“等待状态”的数据模型,要处理用户输入进来之后如何跟原执行流合并,还要保证并发场景下状态读写不冲突。用 LangGraph 的时候,一个interrupt()就帮你把这些全部处理掉了。
interrupt的执行原理是:当图执行到interrupt()调用时,LangGraph 会抛出一个特殊信号,图执行器捕获这个信号后,把当前状态和执行位置写入 Checkpoint,然后优雅地停止执行,把中断信息返回给调用方。这时候 Agent 的“执行栈”就被挂起来了,数据库里留下了足够的恢复信息。之后你用Command(resume=...)再次调用同一个thread_id,图会从上次中断的位置精确恢复,interrupt()的返回值就是你传入的 resume 数据。
用一句话概括:interrupt 就是给 Agent 的执行流打了一个断电续传断点,而不是简单地写一条记录。
3.2 手写一个带人工审批的 Agent
下面用订单退款审批这个最常见也最实用的场景来演示。流程设计是:用户发起退款申请 -> Agent 查询订单和客户信息 -> Agent 向审批人发起询问 -> 审批人给出“同意”或“拒绝” -> Agent 执行对应后续动作。
from langgraph.graph import StateGraph, START, END from langgraph.types import interrupt, Command from langgraph.checkpoint.postgres import PostgresSaver from typing import TypedDict class RefundState(TypedDict): order_id: str amount: float customer_name: str decision: str def query_order(state: RefundState): order_info = order_service.get_by_id(state["order_id"]) return {"customer_name": order_info["customer_name"], "amount": order_info["amount"]} def human_approval(state: RefundState): decision = interrupt({ "type": "refund_approval", "order_id": state["order_id"], "amount": state["amount"], "customer": state["customer_name"], }) return {"decision": decision} def process_decision(state: RefundState): if state["decision"] == "approve": order_service.mark_refunded(state["order_id"]) return {"result": "退款已批准并执行"} else: order_service.mark_rejected(state["order_id"]) return {"result": "退款被拒绝"} graph = StateGraph(RefundState) graph.add_node("query_order", query_order) graph.add_node("human_approval", human_approval) graph.add_node("process_decision", process_decision) graph.add_edge(START, "query_order") graph.add_edge("query_order", "human_approval") graph.add_edge("human_approval", "process_decision") graph.add_edge("process_decision", END) app = graph.compile(checkpointer=checkpointer)这段代码里有几个地方值得单独拎出来说。
human_approval节点做的事情非常纯粹:调用interrupt()把需要展示给用户的信息抛出去,然后等待返回值。返回值会直接作为interrupt()的调用结果被赋值给decision。这里你会发现节点的写法跟普通函数几乎一样,完全不需要手动去轮询数据库或者处理 sleep 逻辑。这就是框架帮你把“挂起/恢复”封装好了的体现。
3.3 触发中断、恢复执行、状态回放
第一次调用的时候,图会执行到human_approval节点然后停下来:
config = {"configurable": {"thread_id": "refund-order-2025-001"}} result = app.invoke({"order_id": "PO-2025-001", "amount": 199.0}, config) # 执行结果是中断信息 print(result)调用方收到的是一个中断结果,里面带着我们传给interrupt()的那份数据,比如这笔退款涉及的订单号、金额和客户姓名。这时候前端可以把这个信息渲染成审批弹窗。要注意,这次调用之后,图并没有执行完毕,它只是“停在半路”,状态已经落到了 PostgreSQL 里。
然后审批人点击“同意”,后端直接恢复执行:
result = app.invoke(Command(resume="approve"), config) print(result) # 输出: {'order_id': 'PO-2025-001', 'amount': 199.0, 'customer_name': '陈*', 'decision': 'approve', 'result': '退款已批准并执行'}这里有个细节:第二次调用的时候,传的还是同一个config,也就是同一个thread_id。LangGraph 会去 Checkpoint 里找上次执行的位置,而不是重新开始。如果你换了thread_id,它就是一个全新的任务,跟之前那个没有任何关系。
调试状态历史的时候,我特别喜欢用get_state_history,它像看 git log 一样清爽:
history = list(app.get_state_history(config)) for snapshot in history: print(snapshot.next, snapshot.values)这个命令会把当前线程从开始到现在的所有 Checkpoint 都列出来,你可以看到每一步执行到了哪个节点、状态值是什么。排查“为什么 Agent 走了一条奇怪的路径”的时候,这个命令基本是必用的。
3.4 用几种方式处理中断分支
实际业务里,审批往往不是简单的二选一。LangGraph 的interrupt()有几种典型用法,我总结一下:
单层中断:就是上面演示的场景,一次等待一个外部输入。适用于流程里的单个审批点或者单次确认。
多次中断:在一条执行链路上多个节点都调用interrupt(),每个节点的 resume 值独立。适用于“节点A等待审批1 -> 执行动作 -> 节点B等待审批2”这种长流程。
恢复值校验:resume 的值可以是你自定义的任意对象,不仅限于字符串。比如传一个 JSON 对象{"approved": True, "note": "同意但要求提供发票"}。如果外部系统返回格式不合法,你可以在human_approval节点里做校验,校验失败再抛一次interrupt()。这样就能实现“重新询问”的效果。
用Command(resume=...)以外的恢复方式:LangGraph 还支持在恢复时直接修改部分状态再继续执行,这在一些需要“改写历史”的场景里很有用。比如用户过来说“上一句话我打错了”,你可以先更新状态里的用户消息,再告诉图从指定节点继续。
4. 完整实操:PostgreSQL 配置与 AG-UI 事件桥接
4.1 PostgreSQL 版本选择与建库建议
操作系统层面我推荐直接用 Docker 起一个 PostgreSQL 16,镜像稳定且内存占用可控。如果公司有现成的 RDS 实例,直接连也行,但为了不干扰线上库,建议单独建一个数据库:
CREATE DATABASE agent_runtime; CREATE USER agent_user WITH PASSWORD 'change_me'; GRANT ALL PRIVILEGES ON DATABASE agent_runtime TO agent_user;这里补充一个版本建议:PostgreSQL 16 是目前兼容性和稳定性都很好平衡的版本,17 也可以用但没必要追新。本地开发图省事可以直接用安装包自带的pgAdmin连接,或者使用便携版启动服务,生产环境则用 Docker 或者云数据库托管,省去很多运维麻烦。
4.2 给生产环境用的初始化脚本
如果你不想在 Python 进程里执行setup(),也可以直接在数据库里手动初始化,效果完全一样。PostgresSaver 的setup()方法执行的 DDL 可以提取出来,放到数据库迁移脚本里统一管理,这样发布流程更规范。大致包含这几件事:
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, type TEXT NOT NULL, metadata JSONB, PRIMARY KEY (thread_id, checkpoint_ns, checkpoint_id) ); CREATE TABLE IF NOT EXISTS checkpoint_blobs ( thread_id TEXT NOT NULL, checkpoint_ns TEXT NOT NULL DEFAULT '', checkpoint_id TEXT NOT NULL, blob BYTEA NOT NULL, PRIMARY KEY (thread_id, checkpoint_ns, checkpoint_id) ); CREATE TABLE IF NOT EXISTS checkpoint_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 NOT NULL, blob BYTEA NOT NULL, PRIMARY KEY (thread_id, checkpoint_ns, checkpoint_id, task_id, idx) );实际线上执行时,我会把这三张表放到agent_runtime数据库的publicschema 里,目录隔离靠thread_id完成。如果是多租户系统,建议再加一列tenant_id做数据隔离,或者直接用不同的 schema 做隔离,这是后话。
4.3 连接池初始化与常驻服务结合
生产环境的 Agent 服务通常是一个常驻的 FastAPI 或 Celery Worker 进程。连接池要在进程启动时初始化好,而不是每次请求进来才创建:
from psycopg_pool import ConnectionPool from langgraph.checkpoint.postgres import PostgresSaver pool = ConnectionPool( conninfo="postgresql://agent_user:change_me@pg-host:5432/agent_runtime", min_size=4, max_size=20, open=False, ) pool.open() checkpointer = PostgresSaver(pool) checkpointer.setup()min_size和max_size的配置要看你的流量模型。Agent 这种任务的特点是长连接比较多、阶段性地写 Checkpoint,不会像 Web 请求那样高频率短连接。我压测下来的经验值是min_size=4、max_size=20,配合 PostgreSQL 默认的max_connections=100,已经能支撑几十路并发 Agent 任务了。如果并发更高,优先扩容数据库连接数,而不是加大 Python 侧连接池。
4.4 AG-UI 在整套架构里的定位
前面讲的都是 LangGraph 内部的事情:图怎么调度、状态怎么存储、中断怎么恢复。但一个完整的 Runtime 不能只有内部机制,它还得面向外部系统有标准接口。这就是 AG-UI 登场的时机。
AG-UI 是 LangChain 团队提出的 Agent Gateway 接口规范,核心目标是标准化 Agent 与前端、桌面端以及其他后端服务之间的通信协议。它定义了一套 HTTP 端点:创建会话、订阅事件流、发送用户输入、获取运行状态等,还定义了一套结构化的事件类型,比如agent消息、tool_call工具调用、interrupt中断请求、error错误、done完成等。
在架构上理解 AG-UI 和 LangGraph 的关系很简单:LangGraph 是运行时引擎,AG-UI 是面向外部的契约。LangGraph 负责“怎么跑”,AG-UI 负责“怎么对外说”。前端不用关心你内部是 LangGraph 还是别的框架,它只需要遵循 AG-UI 规范去拉事件、投输入就行了。
4.5 把 LangGraph 的 state 翻译成 AG-UI 事件流
实际接入的时候,需要做一个轻量的适配层:监听 LangGraph 的stream事件,翻译成 AG-UI 协议的事件格式,再通过 WebSocket 或者 SSE 推给前端。
from fastapi import FastAPI, WebSocket from langgraph.types import StreamMode app = FastAPI() @app.websocket("/session/{session_id}/events") async def session_events(websocket: WebSocket, session_id: str): await websocket.accept() config = {"configurable": {"thread_id": session_id}} async for event in app_graph.astream( {"order_id": "PO-2025-001", "amount": 199.0}, config, stream_mode=StreamMode.UP, ): # 翻译成 AG-UI 兼容的客户端事件 agui_event = { "type": "agent", "session_id": session_id, "timestamp": datetime.now(timezone.utc).isoformat(), "payload": event, } await websocket.send_json(agui_event) # 遇到 interrupt 时额外发一个 interrupt 事件 if event.get("__interrupt__"): await websocket.send_json({ "type": "interrupt", "session_id": session_id, "payload": event["__interrupt__"][0]["value"], })这一段是典型的桥接代码:底层跑的是 LangGraph 的流式事件,上层吐出来的是 AG-UI 风格的标准结构。前端订阅到interrupt事件之后,就知道要把审批弹窗渲染给用户;用户点击按钮,前端往/session/{session_id}/input发一条消息,后端接到后转为Command(resume=...)继续执行。整个链条就通起来了。
这里我特别推荐给事件都加上统一的timestamp,别小看这个字段。多个 Agent 实例并发跑的时候,前端需要按时间线渲染事件,没有统一时钟的时间戳,日志排查和前端展示都会乱套。
4.6 AG-UI 事件格式的速查
为了便于大家在实现时对照,我把 AG-UI 里最常用的事件字段整理了一下:
| 事件类型 | 语义 | 关键字段 |
|---|---|---|
heartbeat | 保活信号 | session_id,timestamp |
agent | Agent 的文字输出或状态消息 | message,source,timestamp |
tool_call | 工具调用开始 | name,args,call_id |
tool_result | 工具调用返回 | name,result,call_id |
interrupt | 需要用户介入 | payload,session_id |
done | 任务完成 | session_id,timestamp |
error | 出错信息 | error,session_id |
具体的字段名和枚举可以以 AG-UI 规范的当前版本为准,但整体结构就是上面这个路子。前端拿到这套事件流,就能展示“Agent 正在调用订单查询工具”“Agent 正在等待审批”这样的实时状态,而不是面对一个永远在转圈的加载动画。
5. 常见问题与排查技巧实录
5.1 thread_id 用错导致状态串线或恢复失败
这是我见过最多的坑。很多同学会把thread_id跟用户 ID 混淆,于是所有用户的任务都共用一个thread_id,结果就是 A 用户的审批弹窗弹给了 B 用户,或者 B 用户恢复的是 A 用户的流程。
正确做法是:每个独立的业务任务/会话都有一个唯一的thread_id,可以用refund-order-2025-001这种业务号,也可以用 UUID。同一个业务任务的多次调用必须复用同一个thread_id,不同任务必须用不同的thread_id。
判断依据很简单:thread_id约等于对话线程 ID,不是用户 ID。一个用户可以同时有多个线程,一个线程也可以跨越多个用户流转。
5.2 恢复执行时返回值一直是 None
有朋友反馈app.invoke(Command(resume="approve"), config)回来了,但decision变量是空的。这个多半是因为Command的路径不对。Command(resume=...)只能解决“上次中断点处的恢复”,如果你在恢复的会话里重新传入了新的输入参数,或者删除了中途的某个节点,图的执行路径就变了,interrupt()可能根本没有被触发。
另外,Command要放在invoke的参数位置,而不是跟初始输入混在一起:
# 正确写法 result = app.invoke(Command(resume="approve"), config) # 错误写法(传了多个输入) # result = app.invoke({"input": "xxx"}, Command(resume="approve"), config)如果还不生效,优先用app.get_state(config)看当前停在哪里:
state_snapshot = app.get_state(config) print(state_snapshot.next)next字段会告诉你图下一步准备执行哪个节点,如果它显示空,说明图并没有停在等待恢复的位置,排查方向就要调整。
5.3 并发场景下 checkpoint 写入抛锁冲突
PostgreSQL 的并发能力远强于 SQLite,但也不是没有坑。最常见的是多个并行节点同时写同一个thread_id下的多个task_id时,由于checkpoint_writes表的主键包含(thread_id, checkpoint_ns, checkpoint_id, task_id, idx),如果并行分支里出现了task_id冲突,就会出现主键冲突报错。
LangGraph 在设计上已经尽量避免这个问题,并行分支的 task_id 是引擎内部生成的,正常情况不会重复。但如果你的图里有手动指定的任务 ID,或者多个进程同时去跑同一个thread_id(比如前端重复点击提交),就会产生写入竞争。
解决思路有两个:一是业务侧保证同一个thread_id同一时刻只有一个实例在跑,用分布式锁或者状态机来控制;二是数据库层给三张表补上适当的索引,避免大表扫描。我在项目的 PostgreSQL 里主要加了这几条索引:
CREATE INDEX idx_checkpoints_thread_id ON checkpoints (thread_id); CREATE INDEX idx_checkpoint_blobs_thread_id ON checkpoint_blobs (thread_id); CREATE INDEX idx_checkpoint_writes_thread_id ON checkpoint_writes (thread_id);5.4 checkpoint 表无限膨胀怎么清理
运行一段时间后,checkpoint_blobs的体积会涨得很快,因为每条状态快照都被完整序列化保存。LangGraph 的get_state_history需要这些历史,所以不能简单粗暴地清空,但要定期归档。
我的做法是:给长时间不活跃的线程做冷数据迁移。比如 30 天没有更新的thread_id,把它相关的三张表记录导出到归档库,然后从主库里删除。这里有一个技巧:检查点之间是通过parent_checkpoint_id关联的,删除时要按thread_id维度整体删,不能只删某一条,否则会把父子链路打断,后续恢复会直接失败。
-- 以 thread 维度清理 30 天前的历史检查点(示例) DELETE FROM checkpoint_writes WHERE thread_id IN ( SELECT thread_id FROM checkpoints WHERE checkpoint_id IN ( SELECT checkpoint_id FROM checkpoints WHERE updated_at < NOW() - INTERVAL '30 days' ) );具体清理策略要按业务容忍度来定,但记住一个原则:正在活跃中的线程,不要动它的任何一条 checkpoint 记录,否则当前运行状态会丢失。
5.5 AG-UI 事件乱序和丢失的问题
AG-UI 走 WebSocket 推事件时,我遇到过事件乱序,特别是 Agent 快速连续调用多个工具时,前端渲染出来的顺序经常前后颠倒。原因是 LangGraph 的astream有多个通道,不同通道的事件到达顺序并不保证。
我的解决办法是,在桥接层不直接转发事件,而是先落一个顺序缓冲区:每条事件打上单调递增的序号再推给前端。前端按序号渲染,就能保证顺序。序号直接用 Python 进程内的itertools.count()生成即可,不需要全局唯一,只需要保证同一个session_id内递增就行。
from itertools import count seq = count(1) # 桥接层 event["seq"] = next(seq) await websocket.send_json(event)如果事件量大且对前端实时性要求高,可以在网关层用 Redis 做缓冲队列,让前端从队列里顺序消费。但大多数场景下,进程内计数器加 WebSocket 的顺序投递已经够用。别在这块过度设计。
6. 结合实战给新人的三条建议
第一个建议:先把手写 Loop 里“状态”和“控制流”分离。哪怕还没引入 LangGraph,先把状态存盘再想循环逻辑,后面的迁移成本会低很多。我当初就是从手工维护messages列表起步的,每次想加一个分支判断,都觉得自己在写面条代码。
第二个建议:生产环境别用 InMemory Checkpoint,哪怕是 demo 撑场面。我一开始嫌 PostgreSQL 配置麻烦,用 MemorySaver 顶了一个星期,结果一次部署重启,所有进行中的用户会话全部消失,客服主管直接来找我喝茶。从那以后,我的项目里只要是能跑的 LangGraph,一律上持久化 Checkpoint,数据库就当是日志存储,这点成本跟出事后的损失完全不成比例。
第三个建议:AG-UI 的接入不要太晚启动。如果后端 Runtime 已经写得很深,前端还停留在自定义事件协议,两边联调会非常痛苦。最好在架构设计阶段就确定一套对外事件格式,LangGraph 内部怎么改都不用动前端。
最后再分享一个小技巧:调试中断恢复时,用get_state_history配合get_state一起看,前者看完整时间线,后者看当前状态。遇到“从中间恢复后状态不对”的玄学问题,十有八九是历史链路里某个节点的状态被意外覆盖了,这时候把两张表的输出放在一起对比,问题基本就能定位。
这套架构跑通之后,我再也不担心进程崩溃导致业务中断了。毕竟 Agent 的价值在于稳定地完成任务,而不只是跑得飞快。