☰
LangGraph+PostgreSQL构建可中断续跑的大模型Agent Runtime
2026/9/29 18:22:10 网站建设 项目流程

1. 项目概述:为什么“中断恢复”成了大模型应用落地的生死线

我第一次在生产环境里跑一个需要调用外部API、等待用户输入、再继续推理的Agent流程时,服务器突然断电了。整个流程卡在“等待用户上传合同PDF”的环节,重启服务后,系统完全不记得刚才干到哪一步——它直接从头开始,又问了一遍“请上传合同”,而用户早已不耐烦地关掉了页面。这种体验不是Bug,是架构缺陷。LangGraph本身不保存状态,它的StateGraph只在内存里跑一次;LangChain的RunnableWithMessageHistory也只管对话轮次,不管业务逻辑的断点。真正让这个项目立住脚的,不是“用了LangGraph”,而是我们把Runtime的生命周期从“一次请求-一次响应”拉长到了“一次任务-多次交互-可中断续跑”。核心关键词就三个:LangGraph负责定义节点与边的拓扑逻辑,PostgreSQL作为Checkpoint存储引擎持久化每一步的中间状态,AG-UI则把“中断点”变成用户可感知、可操作的界面按钮。这不是炫技,是解决真实场景里最痛的三件事:用户中途离开不丢进度、后台任务失败后不重跑全链路、运维人员能随时查看某个订单当前卡在哪一环。我试过用Redis存Checkpoints,但数据过期策略和事务一致性太难把控;也试过SQLite,单机文件锁在并发写入时直接阻塞。最终选PostgreSQL,不是因为它多酷,而是它原生支持INSERT ... ON CONFLICT DO UPDATE(upsert)、行级锁、WAL日志保证崩溃恢复,还有pg_stat_activity能实时查出哪个会话卡在了UPDATE checkpoints SET state = ...上——这些细节,才是“可恢复Runtime”能稳住的关键。

2. 整体架构设计:三层解耦,各司其职

2.1 为什么必须分层?手写Loop的教训太深刻

早期我们用纯Python手写Loop,代码像这样:

while True: state = load_state(task_id) if state["status"] == "waiting_for_upload": handle_upload(state) elif state["status"] == "processing_pdf": extract_text(state) elif state["status"] == "awaiting_approval": send_to_manager(state) # ... 二十多个elif save_state(state) if state["status"] == "completed": break

表面看很清晰,实则埋了三颗雷:第一,save_state()如果写一半崩溃,状态就脏了;第二,新增一个“二次审核”节点,得改所有elif分支,还容易漏掉save_state()调用;第三,UI要显示当前状态,得去解析state["status"]字符串,前端硬编码一堆if (status === 'awaiting_approval')。后来我们意识到:状态迁移逻辑不该由业务代码硬编码,而该由图结构驱动;状态存储不该和业务逻辑混在一起;状态可视化不该依赖字符串匹配。于是拆成三层:LangGraph层只管“什么节点能连到什么节点”,PostgreSQL层只管“把当前节点ID、输入数据、输出数据、时间戳原样存下来”,AG-UI层只管“读取最新Checkpoint,渲染对应UI组件”。这三层之间没有直接调用,全靠task_id这个唯一键串联。比如当用户点击“同意合同”按钮,AG-UI不调用任何Python函数,只发一个HTTP PATCH到/api/tasks/{task_id}/resume,后端收到后,从PostgreSQL里查出task_id对应的最新Checkpoint,提取其中的next_node字段,再调用LangGraph的app.invoke()传入该节点名——整个过程,LangGraph甚至不知道自己被谁调用。

2.2 LangGraph层:用StateGraph定义“合法路径”,而非手写分支

LangGraph的核心价值,在于把“流程控制权”从开发者手里交还给图结构。我们没用MessageGraph,因为消息流不适合我们的业务——合同审批不是聊天,而是有明确输入输出契约的步骤链。我们定义了一个强类型State:

from typing import TypedDict, Optional, List from langgraph.graph import StateGraph class ContractState(TypedDict): task_id: str uploaded_pdf: Optional[str] # S3 URL extracted_text: Optional[str] approval_status: Optional[str] # "pending", "approved", "rejected" manager_notes: Optional[str] current_step: str # "upload", "extract", "review", "sign" last_updated: str

然后构建图:

def upload_node(state: ContractState) -> ContractState: # 调用FastAPI上传接口,返回S3 URL s3_url = call_upload_api(state["task_id"]) return {"uploaded_pdf": s3_url, "current_step": "extract"} def extract_node(state: ContractState) -> ContractState: # 调用PDF解析服务 text = call_pdf_parser(state["uploaded_pdf"]) return {"extracted_text": text, "current_step": "review"} def review_node(state: ContractState) -> ContractState: # 发送邮件给经理,返回待办ID todo_id = send_review_request(state["task_id"], state["extracted_text"]) return {"current_step": "awaiting_approval", "review_todo_id": todo_id} # 构建图:每个节点返回的字典会自动merge进state builder = StateGraph(ContractState) builder.add_node("upload", upload_node) builder.add_node("extract", extract_node) builder.add_node("review", review_node) builder.add_node("sign", sign_node) # 定义边:从upload节点出发,只有成功才走到extract builder.add_edge("upload", "extract") builder.add_edge("extract", "review") builder.add_conditional_edges( "review", lambda state: state["approval_status"], { "approved": "sign", "rejected": END, "pending": "review" # 等待用户操作,不自动推进 } )

关键点在于add_conditional_edges:它让review节点的输出决定下一步,而不是写死if state["approval_status"] == "approved"。这样,当产品说“加个法务复核环节”,我们只需新增一个legal_review_node,改一行{"approved": "legal_review"},不用碰任何if/else。LangGraph的app.invoke()方法接收state和config,其中config里的"thread_id"就是task_id,它会自动从Checkpoint存储中加载上次中断的状态。这里我们没用LangGraph内置的MemorySaver,因为它的内存实现无法跨进程,而我们的Web服务是Gunicorn多Worker部署——必须用外部存储。

2.3 PostgreSQL Checkpoint层:不只是存JSON,而是建业务表

LangGraph官方文档说“Checkpoint可以存任何地方”,但真用起来,你会发现存JSON字符串远远不够。我们最初按官方示例,把整个state序列化成JSON存在一张checkpoints表里:

CREATE TABLE checkpoints ( thread_id TEXT PRIMARY KEY, checkpoint JSONB NOT NULL, created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW() );

两周后就遇到问题:运营要查“今天有多少合同卡在法务复核环节”,SQL得写成SELECT COUNT(*) FROM checkpoints WHERE checkpoint->>'current_step' = 'legal_review',但PostgreSQL对JSONB字段的->>操作符无法走索引,全表扫描。更糟的是,当state里嵌套了几十层字典,checkpoint字段动辄5MB,单条记录写入慢,备份也吃力。我们重构为四张表:

-- 主任务表:存元信息,高频查询 CREATE TABLE tasks ( id SERIAL PRIMARY KEY, task_id TEXT UNIQUE NOT NULL, status VARCHAR(20) NOT NULL CHECK (status IN ('running', 'paused', 'completed', 'failed')), created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(), updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW() ); -- Checkpoint主表:只存关键字段,强制索引 CREATE TABLE checkpoints ( id SERIAL PRIMARY KEY, task_id TEXT NOT NULL REFERENCES tasks(task_id), node_name VARCHAR(50) NOT NULL, -- "upload", "extract" input_data JSONB, -- 输入参数,如{"file_id": "abc123"} output_data JSONB, -- 输出结果,如{"s3_url": "https://..."} created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(), INDEX idx_task_node ON checkpoints(task_id, node_name), INDEX idx_node_created ON checkpoints(node_name, created_at) ); -- 事件日志表:审计用,不参与业务逻辑 CREATE TABLE task_events ( id SERIAL PRIMARY KEY, task_id TEXT NOT NULL, event_type VARCHAR(30) NOT NULL, -- "node_started", "node_completed", "user_resumed" details JSONB, created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW() ); -- 用户操作表:存用户点击行为,用于分析 CREATE TABLE user_actions ( id SERIAL PRIMARY KEY, task_id TEXT NOT NULL, action VARCHAR(30) NOT NULL, -- "approve", "reject", "upload_file" user_id TEXT, created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW() );

现在,运营查“卡在法务复核的合同数”,SQL变成:

SELECT COUNT(*) FROM tasks t JOIN checkpoints c ON t.task_id = c.task_id WHERE t.status = 'running' AND c.node_name = 'legal_review';

idx_task_node索引让这个查询毫秒级返回。更重要的是,input_data和output_data分开存,避免每次更新都重写整个大JSON。比如用户上传文件后,我们只INSERT一条node_name='upload'的记录;法务审批后,再INSERT一条node_name='legal_review'的记录。LangGraph的get_checkpoint()方法被我们重写,它不再读单条JSON,而是按task_id查checkpoints表里created_at最大的那条记录,output_data字段就是下一次invoke()的输入state。这种设计牺牲了一点“状态快照”的原子性,但换来了可运维性——DBA能随时SELECT * FROM checkpoints WHERE task_id = 'xxx' ORDER BY created_at DESC LIMIT 5看到整个执行轨迹。

2.4 AG-UI层:把“中断点”变成用户可操作的按钮

AG-UI不是简单的前端框架,它是连接用户意图和Runtime状态的翻译器。它的核心逻辑就一条:UI组件的状态,必须100%由当前Checkpoint的node_name和output_data决定,不能有任何本地状态。比如“合同上传”组件,它的React代码长这样:

// UploadStep.tsx const UploadStep = ({ taskId }: { taskId: string }) => { const [isUploading, setIsUploading] = useState(false); // 关键:useEffect只依赖taskId,每次taskId变,就重新fetch useEffect(() => { const fetchState = async () => { const res = await fetch(`/api/checkpoints/latest?task_id=${taskId}`); const checkpoint = await res.json(); // 如果当前节点是upload,说明还没上传;如果是extract,说明已上传 if (checkpoint.node_name === 'upload') { // 显示上传按钮 } else if (checkpoint.node_name === 'extract') { // 显示“已上传,正在解析”提示 } }; fetchState(); }, [taskId]); const handleUpload = async (file: File) => { setIsUploading(true); // 调用后端上传接口,后端会INSERT一条node_name='upload'的checkpoint await fetch('/api/upload', { method: 'POST', body: file }); // 上传成功后,强制刷新checkpoint window.location.reload(); // 简单粗暴,确保UI同步 }; return <div>{/* 渲染上传UI */}</div>; };

这里有两个反直觉的设计:第一,我们没用WebSocket实时推送状态变更,因为用户可能关掉页面几小时,再回来时需要的是“最终一致”,不是“实时一致”;第二,window.location.reload()看似暴力,实则可靠——它规避了前端状态管理的所有坑。AG-UI的路由规则是:/task/:taskId根据task_id查tasks表的status,如果是paused或running,就加载对应node_name的组件;如果是completed,就跳转到成功页。当用户点击“同意合同”,AG-UI不调用sign_node,只发PATCH /api/tasks/{taskId}/resume,后端收到后,查checkpoints表里node_name='review'的记录,提取output_data里的review_todo_id,再调用法务系统API完成审批,最后INSERT一条node_name='sign'的记录。整个过程,AG-UI就像一个哑终端,只负责展示和触发,不参与任何业务决策。

3. 核心实现细节:从零搭建可恢复Runtime的七步实操

3.1 第一步:初始化PostgreSQL Checkpoint后端(含事务安全)

LangGraph的Checkpoint接口要求实现AsyncCheckpointSaver抽象类。我们写的PostgresSaver必须解决两个致命问题:并发写入冲突和崩溃后状态不一致。先看基础骨架:

from langgraph.checkpoint.postgres import AsyncPostgresSaver from langgraph.checkpoint.base import Checkpoint, CheckpointMetadata, CheckpointTuple class PostgresSaver(AsyncPostgresSaver): def __init__(self, conn_string: str): super().__init__(conn_string) # 初始化连接池,设置最大连接数为20,避免耗尽DB连接 self.pool = AsyncConnectionPool( conn_string, min_size=5, max_size=20, open=False ) async def aget_tuple(self, config: RunnableConfig) -> Optional[CheckpointTuple]: # 重写aget_tuple:不查langgraph默认表,查我们自己的checkpoints表 async with self.pool.acquire() as conn: # 按task_id查最新checkpoint row = await conn.fetchrow( "SELECT * FROM checkpoints WHERE task_id = $1 ORDER BY created_at DESC LIMIT 1", config["configurable"]["thread_id"] ) if not row: return None # 构造CheckpointTuple,注意:state是output_data,不是整条记录 checkpoint = Checkpoint( ts=row["created_at"].isoformat(), channel_values={}, pending_sends=[], version="1", metadata={"node_name": row["node_name"]} ) # 这里关键:state必须是dict,且包含所有需要的字段 # 我们把output_data反序列化,再merge进初始state模板 initial_state = { "task_id": row["task_id"], "current_step": row["node_name"], "last_updated": row["created_at"].isoformat() } if row["output_data"]: initial_state.update(row["output_data"]) return CheckpointTuple( config=config, checkpoint=checkpoint, metadata=CheckpointMetadata({"source": "postgres"}), pending_tasks=[], parent_config=None )

但光这样还不够。当两个请求同时想更新同一个task_id的状态(比如用户双击“同意”按钮),INSERT会冲突。我们用PostgreSQL的ON CONFLICT DO UPDATE解决:

async def aput(self, config: RunnableConfig, checkpoint: Checkpoint, metadata: CheckpointMetadata) -> RunnableConfig: async with self.pool.acquire() as conn: # 先查当前最新checkpoint的id,用于乐观锁 latest = await conn.fetchrow( "SELECT id FROM checkpoints WHERE task_id = $1 ORDER BY created_at DESC LIMIT 1", config["configurable"]["thread_id"] ) # INSERT新记录,如果task_id已存在(即有未完成的),则更新 # 注意:这里用upsert,但只更新output_data和created_at,不覆盖node_name await conn.execute( """ INSERT INTO checkpoints (task_id, node_name, input_data, output_data, created_at) VALUES ($1, $2, $3, $4, $5) ON CONFLICT (task_id) DO UPDATE SET output_data = EXCLUDED.output_data, created_at = EXCLUDED.created_at """, config["configurable"]["thread_id"], checkpoint["metadata"]["node_name"], # 从checkpoint里取节点名 json.dumps(checkpoint.get("input_data", {})), json.dumps(checkpoint.get("output_data", {})), datetime.now(timezone.utc) ) # 同时更新tasks表的status await conn.execute( "UPDATE tasks SET status = $1, updated_at = $2 WHERE task_id = $3", "running" if checkpoint["metadata"]["node_name"] != "completed" else "completed", datetime.now(timezone.utc), config["configurable"]["thread_id"] ) return config

这里有个精妙点:ON CONFLICT (task_id)不是按主键冲突,而是我们给checkpoints表加了UNIQUE (task_id)约束。这意味着每个task_id最多只有一条“活跃”记录,新状态总是覆盖旧状态。虽然丢失了历史轨迹,但换来了强一致性——用户永远看到最新的状态。如果需要审计,查task_events表即可。

3.2 第二步:配置LangGraph App,启用Checkpoint并注入PostgresSaver

LangGraph的StateGraph本身不处理Checkpoint,必须通过CompiledGraph的checkpointer参数注入。我们没用MemorySaver,而是传入自定义的PostgresSaver实例:

from langgraph.graph import StateGraph from langgraph.checkpoint import BaseCheckpointSaver # 初始化PostgresSaver saver = PostgresSaver("postgresql://user:pass@localhost:5432/mydb") # 构建图 builder = StateGraph(ContractState) # ... 添加节点和边(见2.2节) # 编译图,并注入checkpointer app = builder.compile( checkpointer=saver, # 关键:设置interrupt_before=["review", "legal_review"] # 表示在进入review和legal_review节点前暂停,等待用户操作 interrupt_before=["review", "legal_review"] ) # 启动服务时,确保数据库连接池已初始化 @app.on_event("startup") async def startup(): await saver.pool.open()

interrupt_before是LangGraph可恢复性的灵魂。它告诉LangGraph:“当流程即将进入review节点时,别执行它,先停下来,把当前state存到Checkpoint,然后等外部信号”。这个“外部信号”就是AG-UI发来的/resume请求。app.invoke()方法在遇到interrupt时,会抛出GraphInterrupted异常,但我们不捕获它——让异常冒泡到FastAPI的全局异常处理器,在那里我们记录日志并返回202 Accepted,告诉前端“已暂停,等你操作”。

3.3 第三步:实现/resume端点,安全地恢复执行

/api/tasks/{task_id}/resume端点是Runtime的“心脏起搏器”。它必须做三件事:校验用户权限、加载Checkpoint、触发LangGraph执行。代码如下:

from fastapi import APIRouter, HTTPException, Depends from sqlalchemy.ext.asyncio import AsyncSession router = APIRouter() @router.patch("/tasks/{task_id}/resume") async def resume_task( task_id: str, current_user: User = Depends(get_current_user), # JWT鉴权 db: AsyncSession = Depends(get_db) ): # 1. 权限校验:检查用户是否有权操作此task task = await db.execute( select(Task).where(Task.task_id == task_id) ) task = task.scalar_one_or_none() if not task: raise HTTPException(404, "Task not found") # 检查用户是否是任务创建者或管理员 if task.user_id != current_user.id and not current_user.is_admin: raise HTTPException(403, "Forbidden") # 2. 加载最新Checkpoint config = {"configurable": {"thread_id": task_id}} try: # LangGraph的get_state会从PostgresSaver里查 state = await app.aget_state(config) if not state or not state.values: raise HTTPException(400, "No checkpoint found for this task") # 3. 触发执行:app.invoke会从中断点继续 # 注意:必须传入config,否则LangGraph不知道从哪恢复 result = await app.ainvoke( input={}, # 输入为空,因为state里已有所有数据 config=config ) # 4. 更新tasks表状态 await db.execute( update(Task).where(Task.task_id == task_id).values( status="completed" if result.get("current_step") == "completed" else "running", updated_at=func.now() ) ) await db.commit() return {"status": "resumed", "next_step": result.get("current_step")} except Exception as e: # 记录详细错误,包括task_id和state快照 logger.error(f"Resume failed for task {task_id}: {str(e)}") await db.rollback() raise HTTPException(500, "Failed to resume task")

这里的关键是app.ainvoke(input={}, config=config)。input={}表示不提供新输入,LangGraph会自动从Checkpoint里加载上次的state,然后从interrupt_before指定的节点继续执行。比如上次停在review,这次就会执行review_node函数。如果review_node里调用了发送邮件的API,邮件就会立刻发出。

3.4 第四步:AG-UI的动态路由与组件加载机制

AG-UI的前端用React + React Router v6实现。它的路由不是静态的,而是根据task_id动态生成:

// App.tsx function App() { return ( <Router> <Routes> {/* 动态路由:/task/:taskId 匹配任意task_id */} <Route path="/task/:taskId" element={<TaskPage />} /> <Route path="/" element={<HomePage />} /> </Routes> </Router> ); } // TaskPage.tsx:根据task_id加载对应组件 const TaskPage = () => { const { taskId } = useParams(); const [currentNode, setCurrentNode] = useState<string | null>(null); const [isLoading, setIsLoading] = useState(true); useEffect(() => { const loadTask = async () => { try { setIsLoading(true); // 1. 查tasks表获取任务状态 const taskRes = await fetch(`/api/tasks/${taskId}`); const task = await taskRes.json(); if (task.status === 'completed') { navigate(`/task/${taskId}/success`); return; } // 2. 查checkpoints表获取当前节点 const cpRes = await fetch(`/api/checkpoints/latest?task_id=${taskId}`); const cp = await cpRes.json(); // 3. 根据node_name决定渲染哪个组件 setCurrentNode(cp.node_name); } catch (e) { console.error(e); } finally { setIsLoading(false); } }; loadTask(); }, [taskId]); if (isLoading) return <LoadingSpinner />; // 动态导入组件,避免打包体积过大 const Component = dynamic(() => import(`./steps/${currentNode}Step`).then(m => m.default)); return ( <div className="task-container"> <Header taskId={taskId} /> <Component taskId={taskId} /> <Footer /> </div> ); };

dynamic import是关键。它让Webpack把每个Step组件打成独立chunk,用户访问/task/abc123时,只加载uploadStep.js,不会下载signStep.js。steps/目录结构如下:

steps/ ├── uploadStep.tsx // 对应node_name='upload' ├── extractStep.tsx // 对应node_name='extract' ├── reviewStep.tsx // 对应node_name='review' └── legalReviewStep.tsx // 对应node_name='legal_review'

每个Step组件内部,都封装了该节点的专属UI和API调用。比如reviewStep.tsx里有“同意”、“拒绝”两个按钮,点击后调用/api/tasks/{taskId}/resume,而不是直接调用后端业务API。这种设计让UI彻底解耦——产品经理说“把‘同意’按钮改成绿色”,我们只改reviewStep.tsx,不影响LangGraph或PostgreSQL。

3.5 第五步:处理超时与失败的兜底策略

生产环境没有“永远在线”。我们设定了三重超时保护:节点超时、任务超时、Checkpoint清理。首先,给每个LangGraph节点加超时:

import asyncio async def upload_node(state: ContractState) -> ContractState: try: # 设置10秒超时 s3_url = await asyncio.wait_for( call_upload_api(state["task_id"]), timeout=10.0 ) return {"uploaded_pdf": s3_url, "current_step": "extract"} except asyncio.TimeoutError: # 超时后,写入失败事件,并标记任务为failed await log_failure(state["task_id"], "upload_timeout") raise NodeFailedError("Upload timed out")

其次,用Celery Beat定时扫描“卡住”的任务:

# celery_tasks.py from celery import Celery app = Celery('tasks', broker='redis://localhost:6379') @app.task def check_stuck_tasks(): # 查找30分钟内没有更新的running任务 query = """ SELECT t.task_id, c.node_name FROM tasks t JOIN checkpoints c ON t.task_id = c.task_id WHERE t.status = 'running' AND c.created_at < NOW() - INTERVAL '30 minutes' """ # 执行查询,对每个stuck task发告警,并调用force_resume for row in execute_query(query): send_alert(f"Task {row['task_id']} stuck at {row['node_name']}") force_resume(row['task_id']) # 强制重试

最后,Checkpoint表加分区和TTL。PostgreSQL 12+支持按时间分区:

-- 创建按月分区的checkpoints表 CREATE TABLE checkpoints_2024_09 PARTITION OF checkpoints FOR VALUES FROM ('2024-09-01') TO ('2024-10-01'); -- 自动清理3个月前的分区 CREATE OR REPLACE FUNCTION cleanup_old_checkpoints() RETURNS void AS $$ BEGIN DROP TABLE IF EXISTS checkpoints_2024_06; END; $$ LANGUAGE plpgsql;

3.6 第六步:本地开发环境快速启动(含PostgreSQL一键安装)

开发者最怕“环境搭三天,代码写五分钟”。我们用Docker Compose统一本地环境:

# docker-compose.yml version: '3.8' services: postgres: image: postgres:15 environment: POSTGRES_DB: myapp POSTGRES_USER: user POSTGRES_PASSWORD: pass ports: - "5432:5432" volumes: - postgres_data:/var/lib/postgresql/data healthcheck: test: ["CMD-SHELL", "pg_isready -U user -d myapp"] interval: 30s timeout: 10s retries: 5 web: build: . environment: DATABASE_URL: postgresql://user:pass@postgres:5432/myapp ports: - "8000:8000" depends_on: postgres: condition: service_healthy volumes: postgres_data:

配套一个init-db.sql脚本,放在docker-entrypoint-initdb.d/目录下,容器启动时自动执行:

-- init-db.sql CREATE TABLE IF NOT EXISTS tasks (...); CREATE TABLE IF NOT EXISTS checkpoints (...); CREATE TABLE IF NOT EXISTS task_events (...); CREATE TABLE IF NOT EXISTS user_actions (...); -- 插入测试数据 INSERT INTO tasks (task_id, status) VALUES ('test-001', 'running'); INSERT INTO checkpoints (task_id, node_name, output_data, created_at) VALUES ('test-001', 'upload', '{"s3_url": "https://test.s3/test.pdf"}', NOW());

开发者只需docker-compose up -d,5秒后就能访问http://localhost:8000/task/test-001看到上传界面。AG-UI的.env文件里,VITE_API_BASE_URL=http://localhost:8000,前后端分离,互不干扰。

3.7 第七步:监控与可观测性——让Runtime“看得见”

没有监控的Runtime就像没有仪表盘的飞机。我们在三个层面埋点:LangGraph层、PostgreSQL层、AG-UI层。LangGraph层用app.add_node的回调:

def log_node_start(state: ContractState, config: RunnableConfig): logger.info( f"Node {config['configurable'].get('node_name', 'unknown')} started for task {state['task_id']}", extra={"task_id": state["task_id"], "node": config["configurable"].get("node_name")} ) # 在编译前注册 app.add_node("upload", upload_node, on_start=log_node_start)

PostgreSQL层用pg_stat_statements扩展,查慢查询:

-- 开启扩展 CREATE EXTENSION IF NOT EXISTS pg_stat_statements; -- 查最慢的10个查询 SELECT query, total_time, calls, total_time/calls as avg_time FROM pg_stat_statements WHERE query LIKE '%checkpoints%' ORDER BY total_time DESC LIMIT 10;

AG-UI层用useEffect监听visibilitychange事件,记录用户离开页面的时间:

useEffect(() => { const handleVisibilityChange = () => { if (document.hidden) { // 用户切走了,记录时间 localStorage.setItem(`task_${taskId}_hidden_at`, Date.now().toString()); } else { // 用户切回来了,计算离线时长 const hiddenAt = localStorage.getItem(`task_${taskId}_hidden_at`); if (hiddenAt) { const duration = Date.now() - parseInt(hiddenAt); // 上报到监控系统 reportUserOfflineDuration(taskId, duration); } } }; document.addEventListener('visibilitychange', handleVisibilityChange); return () => document.removeEventListener('visibilitychange', handleVisibilityChange); }, [taskId]);

所有日志都打到ELK栈,用Kibana看Dashboard:一个面板显示“当前运行中任务数”,另一个显示“平均恢复耗时”,第三个显示“各节点失败率”。当review_node失败率突增,我们立刻知道是法务系统挂了,而不是LangGraph有问题。

4. 实战问题排查:那些文档里不会写的坑

4.1 问题1:LangGraph恢复后,state里丢失了非JSON序列化的对象

现象:upload_node返回{"uploaded_pdf": <S3Object>},但S3Object是boto3的类,无法JSON序列化。恢复时state["uploaded_pdf"]变成None。

根因:LangGraph的Checkpoint存储要求state必须是JSON-serializable。S3Object有方法和属性,json.dumps()直接报错。

解决方案:在节点函数里,只存可序列化的数据。upload_node应该返回{"uploaded_pdf": "s3://bucket/key.pdf"},而不是对象本身。所有业务逻辑里需要S3对象的地方,用boto3.client('s3').get_object()按需获取。我们写了个装饰器强制检查:

import json def ensure_json_serializable(func): async def wrapper(*args, **kwargs): result = await func(*args, **kwargs) try: json.dumps(result) except TypeError as e: raise ValueError(f"Node {func.__name__} returned non-serializable state: {e}") return result return wrapper @ensure_json_serializable async def upload_node(state: ContractState) -> ContractState: # ...

提示:这个装饰器必须加在所有节点函数上。我们把它放进CI流水线,用pytest跑单元测试,确保每个节点的返回值都能json.dumps()。

4.2 问题2:PostgreSQL连接池耗尽,所有请求503

现象:高并发时,web服务日志满屏psycopg2.OperationalError: connection limit exceeded。

根因:我们设了max_size=20,但每个FastAPI请求的app.invoke()会开一个数据库连接,而LangGraph内部可能开多个连接(比如aget_state和aput各一个)。20个连接被10个并发请求占满,第11个请求就卡住。

解决方案:用连接池的acquire(timeout=...)加超时,并降级。修改PostgresSaver:

async def aget_tuple(self, config: RunnableConfig) -> Optional[CheckpointTuple]: try: # 设置5秒超时,超时后返回None,让LangGraph用默认state async with self.pool.acquire(timeout=5.0) as conn: # ... 查询逻辑 except PoolTimeoutError: logger.warning(f"Postgres pool timeout for task {config['configurable']['thread_id']}") return None # LangGraph会继续,但state不完整

同时,把max_size从20降到10,因为app.invoke()是异步的,连接可以复用。我们用pg_stat_activity监控:

SELECT usename, application_name, state, COUNT(*) FROM pg_stat_activity GROUP BY usename, application_name, state;

发现application_name为langgraph-checkpoint的连接数稳定在8-10,证明调优成功。

4.3 问题3:AG-UI加载时,查不到最新Checkpoint

现象:用户上传完文件,页面刷新,却还显示“上传中”,而不是“正在解析”。

根因:AG-UI的fetch('/api/checkpoints/latest')和后端INSERT INTO checkpoints不在同一个数据库事务里。INSERT提交后,SELECT才能看到,但网络延迟导致SELECT在INSERT前发出。

解决方案:在/upload端点里,INSERT后立即SELECT返回最新记录,AG-UI用这个结果,而不是再发一次请求:

#

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询