AI Agent 生产落地:基于 LangGraph 与 K8s 的长任务容错编排实战

音乐学院打鼓的时候,乐队主音吉他手如果偶尔抢拍,贝斯手和鼓手还能凭着节奏感强行把拍子拽回来;但如果鼓手的双踩直接卡死,整首曲子在台上下不来台就是一场灾难。

在生产环境落地大模型 Agent(智能体)也是完全一样的道理。

过去绝大多数演示(Demo)级别的 Agent,本质上只是一个单线程的无状态 while 循环或者线性 Prompt 链。但在企业级云原生生产环境中,一个复杂的 Agent 任务(例如自动化代码审计、长篇市场调研、K8s 集群日志根因排查)往往需要连续运行几十分钟,触发几十次 Tool Calling(工具调用)和模型推理。

在这么长的执行链条中,一旦遭遇网络抖动、模型 API 超时限制(Rate Limit)、或者 K8s 节点被抢占导致 Pod 崩溃,传统无状态 Agent 的整个上下文与推理步骤就会彻底蒸发。这意味着昂贵的 Token 算力直接白烧,SLA(服务等级协议)掉到地缝里。

要让 Agent 具备工业级的防爆恢复能力,必须将 Agent 改造为基于 LangGraph 的有状态图(State Graph),并结合 K8s 的 Pod 容错与 SQLite/PostgreSQL 持久化 Checkpoint(检查点) 机制。


状态图与持久化 Checkpoint 容错架构

传统的 Agent 依赖内存中的 Context 变量,而 LangGraph 的核心思想是将 Agent 抽象为一个强类型的有状态图(StateGraph)

每一个推理步骤、工具调用与条件分支,都是图上的一个独立节点(Node);而 Agent 在整个生命周期中流动的数据,统一收拢在全局状态(State)中。

flowchart TD
    UserTask[用户发起长链条 Agent 任务] --> K8sPod[K8s Pod: 运行 LangGraph 节点器]
    
    subgraph LangGraph 有状态图与 Checkpoint 存储
        K8sPod --> NodePlanner[节点 1: Planner 任务拆解]
        NodePlanner -->|原子更新 State| MemorySaver[持久化 Checkpointer: PostgreSQL / SQLite]
        
        MemorySaver -->|保存 Thread_ID + Checkpoint_ID| DBStore[(持久化存储磁盘)]
        
        NodePlanner --> NodeTool[节点 2: Tool Calling 批量执行]
        NodeTool -->|突发网络故障 / Pod 被 Kill| PodCrash[K8s Pod 异常崩溃退出]
        
        PodCrash -->|K8s 重新拉起新 Pod| NewPod[新 Pod: 读取最新 Checkpoint]
        DBStore -->|恢复最近一次成功状态| NewPod
        NewPod --> NodeEvaluator[节点 3: 无缝续跑 Evaluator 校验]
    end
    
    NodeEvaluator --> FinalOutput[输出最终可靠结果]

1. 持久化 Checkpoint 机制的物理过程

当 Agent 在 K8s 节点上执行时,LangGraph 会在每一个节点转移(State Transition)完成后,触发一次原子级 Checkpoint 写入

  • thread_id:唯一标识当前用户的 Task 会话。
  • checkpoint_id:标识当前步骤的递增版本号。
  • 全局状态 Snapshot:将当前的对话历史、工具返回结果与内部变量序列化后落盘。

2. K8s Pod 崩溃后的秒级无缝续跑

如果正在运行 Tool Calling 的 Pod 突然触发了 OOMKilled 或者被 K8s 抢占,K8s 部署控制器(Deployment / StatefulSet)会立即拉起一个新的镜像容器。

新 Pod 启动后,根据传入的 thread_id 向数据库查询最新一条成功记录的 checkpoint_id,直接从中恢复全局 State,跳过前面已经执行成功的节点,继续往下跑


生产级 Python 代码:LangGraph 持久化 Agent 编排器实现

下面是一套可以在生产环境中直接运行的 Python 源码。它使用了 LangGraph 库建立状态图,配合 SQLite 持久化 Checkpointer,并实现了模拟崩溃与恢复机制:

import os
import sqlite3
import logging
from typing import Dict, Any, TypedDict, Annotated, List
from langgraph.graph import StateGraph, END
from langgraph.checkpoint.sqlite import SqliteSaver

logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s")
logger = logging.getLogger("AgentOrchestrator")

# 1. 定义全局状态结构体
class AgentState(TypedDict):
    task_id: str
    input_query: str
    steps_completed: List[str]
    tool_outputs: Dict[str, Any]
    current_status: str
    retry_count: int

# 2. 定义各个节点处理逻辑
def planner_node(state: AgentState) -> Dict[str, Any]:
    logger.info(f"[Node: Planner] 正在为任务 {state['task_id']} 拆解步骤...")
    completed = state.get("steps_completed", [])
    completed.append("planner_done")
    return {
        "steps_completed": completed,
        "current_status": "PLANNING_COMPLETE"
    }

def tool_execution_node(state: AgentState) -> Dict[str, Any]:
    logger.info(f"[Node: ToolExecution] 正在调用外部工具 API...")
    completed = state.get("steps_completed", [])
    
    # 模拟处理逻辑
    tool_data = state.get("tool_outputs", {})
    tool_data["k8s_cluster_status"] = "Node-01 Healthy, Node-02 Ready"
    completed.append("tools_done")

    return {
        "steps_completed": completed,
        "tool_outputs": tool_data,
        "current_status": "TOOLS_COMPLETE"
    }

def evaluator_node(state: AgentState) -> Dict[str, Any]:
    logger.info(f"[Node: Evaluator] 正在对结果进行终态校验...")
    completed = state.get("steps_completed", [])
    completed.append("evaluator_done")
    return {
        "steps_completed": completed,
        "current_status": "FINISHED"
    }

# 3. 构造状态图与条件路由
def build_agent_graph():
    builder = StateGraph(AgentState)

    builder.add_node("planner", planner_node)
    builder.add_node("tool_execution", tool_execution_node)
    builder.add_node("evaluator", evaluator_node)

    builder.set_entry_point("planner")
    builder.add_edge("planner", "tool_execution")
    builder.add_edge("tool_execution", "evaluator")
    builder.add_edge("evaluator", END)

    return builder

if __name__ == "__main__":
    db_path = "/tmp/agent_checkpoints.db"
    
    # 使用 SQLite 作为生产级本地持久化 Checkpoint 数据库
    conn = sqlite3.connect(db_path, check_same_thread=False)
    memory_saver = SqliteSaver(conn)

    builder = build_agent_graph()
    # 编译图并注入 Checkpointer
    agent_app = builder.compile(checkpointer=memory_saver)

    config = {"configurable": {"thread_id": "task_session_9981"}}

    # 4. 初始化启动任务
    initial_state = {
        "task_id": "task_session_9981",
        "input_query": "检查集群故障并自动修复",
        "steps_completed": [],
        "tool_outputs": {},
        "current_status": "STARTING",
        "retry_count": 0
    }

    logger.info("[*] 首次启动 Agent 执行长链条任务...")
    for event in agent_app.stream(initial_state, config):
        logger.info(f"实时状态变迁: {event}")

    logger.info("\n" + "=" * 50)
    logger.info("[*] 模拟 Pod 被被强制杀死重启,直接从数据库读取最新 Checkpoint 恢复...")

    # 5. 模拟新 Pod 启动:只传入相同的 thread_id,无需传入初始参数
    restored_state = agent_app.get_state(config)
    logger.info(f"从磁盘恢复成功!已完成的步骤: {restored_state.values.get('steps_completed')}")
    logger.info(f"终态状态记录: {restored_state.values.get('current_status')}")

    conn.close()
    if os.path.exists(db_path):
        os.remove(db_path)

架构权衡与工程考量(Trade-offs)

在将 LangGraph 容错编排引入云原生环境时,需要客观考量以下工程权衡:

评估维度 传统无状态 Agent (Loop 循环) LangGraph + Checkpoint 容错架构 工程权衡 (Trade-offs)
故障恢复能力 零(崩溃后需全量从头重新执行) 极强(秒级读取磁盘,从崩溃节点无缝续跑) 大幅降低了长 Task 在不稳定网络下的 Token 浪费。
IO 与存储开销 零额外存储 每一步需向 PostgreSQL / SQLite 写入 Snapshot 增加了微量的数据库磁盘写 I/O。
状态设计复杂度 简单 需严格定义 Schema 与原子节点 要求开发者具备严谨的状态机思维,禁止在节点内部混入不可序列化对象。

在大厂高并发云原生架构中,用微量的磁盘 I/O 换取 Agent 执行链路 100% 的确定性与中断恢复能力,是构建企业级 AI 应用最划算的一笔账。


总结

掉拍子的鼓手留不住听众,容易崩溃失忆的 Agent 留不住用户。

通过把大模型 Agent 的运行逻辑收拢为 LangGraph 强类型状态图,结合持久化 Checkpoint 与 K8s 节点容错,才能让 Agent 在面对云原生环境的网络抖动与 Pod 重启时,依然像打了稳固节拍一样死死卡住执行节奏,提供真正工业级的稳定性。


参考资料

Logo

中国智能体开发者社区,聚焦智能体与大模型开发,提供前沿资讯、实用工具链、开源项目及行业案例。通过技术沙龙、开发者大赛等活动,促进经验交流与协作,助力开发者快速构建创新智能应用。

更多推荐