目录

一、LangGraph 的核心组件

二、核心组件应用示例

1. State状态(定义参数)

2. Node普通节点

2.1. Node工具调用节点

3. Edge边的两种写法

3.1. 普通边(固定顺序)

3.2. 条件边(条件判断)

4. 执行工作流

三、LangGraph异常处理(重试、熔断)

1. 重试

1.1. 节点重试

1.2. 工具节点重试

1.3. 熔断机制

四、带人工确认节点的工作流(工作流持久化)

五、构建“反思机制”智能体

它解决什么问题?

和 ReAct 有啥区别?

流程长什么样?

State 里通常放什么?

什么时候用 / 不用?


一、LangGraph 的核心组件

State(状态)[数据传输],所有节点的“共享工作记忆”,存储任务中的全部信息。

        特点:增量更新(节点只返回部分修正)

Node(节点):[最小执行单元],例如调用AI、查询数据库等。

        节点职责通常单一,输入当前 State,并返回 State 的部分更新

Edge(边):[定义工作流执行顺序]普通边(固定顺序执行)、条件边(条件判断)

        支持:循环、回溯、并行、多分支、终止(END)

StateGraph(状态图):[工作流的蓝图] 定义了状态的结构、节点和边的整体布局

        所有节点和边都必须在一个 StateGraph 实例上定义和组合

二、核心组件应用示例

1. State状态(定义参数)

State需要那些字段完全基于业务自定义。

messages: Annotated[list, operator.add] 是常用字段类型,可以累加所有节点需要记录的内容。

messages 存储那些数据,一般有2种场景:

场景一、LLM对话类一般存储,每个节点和大模型的对话记录:

{"role": "user", "content": "用户输入"}
{"role": "assistant", "content": "AI回答"}

# 例如返回当前节点结果【只返回需要更新的字段(字典),LangGraph会自动合并到状态中】
return {
        "messages": [
            {"role": "assistant", "content": “大模型返回的结果文本”}
        ],
        "next_step": "bb"
}

场景二、执行日志 / 流程记录(每个节点记一句干了什么、结果是什么,方便排查):

def aa_node(state: ConstructionState):
    return {
        "messages": [
            {
                "step": "节点名称",
                "action": "节点作用",
                "status": "成功/失败",
                "data": "上一节点参数值,当前节点处理结果"
            }
        ],
        "next_step": "bb"
    }
class ConstructionState(TypedDict):

    # 消息上下文:自动累加追加,不覆盖历史    
    # 作用:所有节点都能看到前面所有步骤的历史
    messages: Annotated[list, operator.add] # Annotated[list, operator.add] 创建一个累加的list参数

    # 【路由控制】控制条件边下一步走向
    # 可选值:"bb", "cc", "dd", "end" 等任意节点名称
    # 判断函数根据该字段决定流程走向
    next_step: str

    # 用户原始输入内容
    input_text: str

    # 工作流最终输出结果(当前节点输出的结果)
    result: str

    # ...

2. Node普通节点

定义三个普通函数方法节点

# ========================
# 节点函数 伪代码
# 规则:函数必须 return 字典格式 ConstructionState → 不返回的字段节点参数值不变。
# ========================

def aa_node(state: ConstructionState):
    """
    第一个节点:入口节点
    作用:处理输入、设置路由、记录消息
    """
    # 1. 读取用户输入(所有state字段都能拿到)
    input_text = state["input_text"]

    # 2. 你的业务逻辑(随便写)
    new_msg = f"[aa节点] 处理用户输入,{input_text} 将问题控制在20字内。"
    result_llm = "调用大模型,返回的接口"

    # 3. 返回要更新的state字段(只返回你要改的!)
    return {
        "messages": [new_msg],        # 自动累加
        "next_step": "bb",            # 控制下一步走向 bb
        "result": result_llm          # 节点结果
    }

def bb_node(state: ConstructionState):
    """
    第二个节点:中间处理节点
    作用:业务处理、加工数据
    """
    # 1. 读取上一步状态
    input_text = state["input_text"]
    prev_result = state["result"]

    # 2. 业务逻辑
    print("=== 执行 bb 节点 ===")
    new_msg = f"[bb节点] 基于aa结果继续处理:{prev_result}"
    new_result = f"{prev_result} → bb节点处理完成"

    # 3. 返回更新
    return {
        "messages": [new_msg],
        "result": new_result
    }

def cc_node(state: ConstructionState):
   # 和以上写法一样

2.1. Node工具调用节点

from langgraph.prebuilt import ToolNode
from langchain_core.tools import tool

# 1. 定义工具
@tool
def get_weather(city: str) -> str:
    """获取天气"""
    return f"{city}天气晴朗"

@tool  
def calculate(expression: str) -> str:
    """计算数学表达式"""
    return str(eval(expression))

tools = [get_weather, calculate]

# 2. ToolNode 是 LangGraph 预置的工具节点(自动执行工具调用,返回最后一条结果)
# 注意:ToolNode 可以调用多个工具,但是不会做循环调用(react)
tool_node = ToolNode(tools)


# 2. 在图中使用(和普通节点使用方式一样)
graph_builder.add_node("tools", tool_node)  # 直接添加,不需要写函数

3. Edge边的两种写法

StateGraph状态图(state状态)、添加node节点、Edge边的完整使用案例。

3.1. 普通边(固定顺序)

开始 → aa → bb  → cc → 结束

def create_construction_workflow() -> StateGraph:
    # 1.定义状态图(定义工作流)【类似找白纸,后续在上边画节点】
    workflow = StateGraph(ConstructionState) # 参数(State 状态)

    # 2.添加节点
    workflow.add_node("aa", aa_node) # (定义节点命名,对应的函数方法)
    workflow.add_node("bb", bb_node)
    workflow.add_node("cc", cc_node)

    # ========================
    # 3.普通边(连接节点),执行顺序aa节点 → bb节点 → cc节点
    # ========================
    workflow.set_entry_point("aa") # 设置入口节点为aa
    workflow.add_edge("aa", "bb") # aa执行完执行bb节点
    workflow.add_edge("cc",  END) # 设置结束节点为bb

    return workflow.compile() # 编译工作流

3.2. 条件边(条件判断)

开始 → aa → 分支判断

                           ├─ go_bb → bb → cc → 结束

                           └─ go_cc → cc → 结束

# 判断函数:aa 执行完,决定去 bb 还是 cc
def judge_aa_to_where(state):
    if state["next_step"]  == "go_bb":
        return "go_bb"   # 对应去 bb
    else:
        return "go_cc"   # 对应去 cc


def create_construction_workflow() -> StateGraph:
    # 1. 定义工作流(状态)
    workflow = StateGraph(ConstructionState)

    # 2. 添加节点
    workflow.add_node("aa", aa_node) # (定义节点命名,对应的函数方法)
    workflow.add_node("bb", bb_node)
    workflow.add_node("cc", cc_node)

    # 3. 设置边入口
    workflow.set_entry_point("aa")

    # ========================
    # 条件边(多分支判断):aa 执行完 → 判断走 bb 还是 cc节点
    # 注意:一次 add_conditional_edges 只能设置【一个节点】的分支跳转。
    # ========================
    workflow.add_conditional_edges(
        "aa",                    # 从 aa节点 出发
        judge_aa_to_where,       # 判断函数
        {                        # 路线表
            "go_bb": "bb",       # 如果judge_aa_to_where函数返回go_bb,就执行bb节点
            "go_cc": "cc"        # 如果judge_aa_to_where函数返回go_cc,就执行cc节点
        }
    )

    # bb 走完 → 去 cc
    workflow.add_edge("bb", "cc")

    # cc 走完 → 结束
    workflow.add_edge("cc", END)

    return workflow.compile() # 编译工作流

4. 执行工作流

    app = create_construction_workflow();# 获取工作流实例

    result = app.invoke(initial_state) # 执行工作流

三、LangGraph异常处理(重试、熔断)

完整代码路径(超时控制、重试机制、异步调用):\week09\4\p34常见异步陷阱及规避

1. 重试

1.1. 节点重试

from langgraph.types import RetryPolicy

# 基本重试策略
basic_policy = RetryPolicy(
    max_attempts=3,       # 最大重试次数(1+3=4次)
    retry_delay=1.0,      # 重试延迟(秒)
    backoff_factor=2.0,   # 退避因子,每次等待时间翻倍(2→4→8秒)
    max_delay=60.0,       # 最大延迟时间
    retry_on=(ValueError,)# 可重试的异常类型
)

...

# retry=basic_policy 加入节点重试规则
builder.add_node("get_user_node", get_user_node, retry=basic_policy)

...

1.2. 工具节点重试

from langgraph.types import RetryPolicy
from langgraph.prebuilt import ToolNode

# 工具节点
tool_node = ToolNode(
    # 配置工具列表
    tools=[api_call_function],
    # 工具执行失败的重试策略(仅对工具本身生效)
    retry_policy=RetryPolicy(
        max_attempts=3,        # 最多重试3次(总共执行1+3=4次)
        retry_delay=2.0,       # 第一次重试前等待2秒
        backoff_factor=2.0,   # 退避因子=2,每次等待时间翻倍(2→4→8秒)
        retry_on=(ConnectionError, TimeoutError)  # 只在这两种异常时重试
    )
)

# 将工具节点添加到流程图中
builder.add_node("api_call", tool_node)

1.3. 熔断机制

# 熔断器三种状态
正常 → 失败重试次数超阈值 → 熔断开启(直接拒绝请求)→ 半开尝试 → 恢复/继续熔断s四

四、带人工确认节点的工作流(工作流持久化)

人工确认节点依赖工作流持久化将每次节点执行后的 State 存入数据库(即工作流持久化),人工确认节点才能中断后恢复执行。

核心用途:

  • 中断恢复 — 在 自定义节点前暂停(例如在confirm节点前暂停),等人确认后继续执行(持久化的核心用途)
  • 进度查询 — 查看当前执行到哪个节点(不持久化也能做,但服务重启后查不到)
  • 历史追溯 — 查看某次请求各节点的执行记录,用于排查问题

D:\dev\python-projects\demo-01-project\project-fs-demo\graph\project_schedule\workflow.py

工作流添加持久化:

# 配置工作流状态存储
# 这里使用 SQLite 数据库存储状态,也可以选择其他存储方式(如 Redis、MongoDB 等)
 memory = SqliteSaver(
     sqlite3.connect("checkpoints.db", check_same_thread=False)
)

# 编译工作流(到达 confirm 节点前自动 return,等待下次 invoke 时从中断处继续)
return workflow.compile(interrupt_before=["confirm"], checkpointer=memory)

D:\dev\python-projects\demo-01-project\project-fs-demo\app.py

# 首次请求接口:提交用户问题,启动工作流
# 需要人工确认 → 返回 need_confirm(含 thread_id);不需要 → 返回 success

# 首次请求接口:提交用户问题,启动工作流
# 需要人工确认 → 返回 need_confirm(含 thread_id);不需要 → 返回 success
@app.post("/api/add/progress", response_model=QueryResponse)
async def handle_add_progress(request: QueryRequest):
    if not request.user_query or not request.user_query.strip():
        raise HTTPException(status_code=400, detail="查询参数不能为空")
    
    try:
        # 初始化State状态值
        initial_state = AddProgressState(
            user_query=request.user_query,
            original_form_info=None,
            uploaded_file=None,
            messages=[]
        )
        

        """
            生成唯一的工作流ID,用于中断恢复(人工确认后继续执行)
        """
        config = {"configurable": {"thread_id": str(uuid.uuid4())}}
        """
            执行工作流
        """
        result = workflow.invoke(initial_state.model_dump(), config=config)
        """
            检查工作流返回结果是否需要人工确认
        """
        state = workflow.get_state(config)
        if state.next == ("confirm",):
            return QueryResponse(
                status="need_confirm",
                result="请确认以下信息是否正确",
                messages=result["confirm_data"],
                thread_id=config["configurable"]["thread_id"]
            )

        return QueryResponse(
            status="success",
            result=result.get("result", "操作完成"),
        )
    except Exception as e:
        raise HTTPException(status_code=500, detail=str(e))

# 二次确认接口:用户确认/拒绝后,从中断处恢复执行
# 用户确认(yes)→ 继续推到飞书;用户拒绝(no)→ 工作流正常结束

# 二次确认接口:用户确认/拒绝后,从中断处恢复执行
# 用户确认(yes)→ 继续推到飞书;用户拒绝(no)→ 工作流正常结束
@app.post("/api/add/progress/confirm", response_model=QueryResponse)
async def handle_confirm(request: ConfirmRequest):
    try:
        config = {"configurable": {"thread_id": request.thread_id}}
        workflow.update_state(config, {"user_confirmed": request.user_confirmed})
        result = workflow.invoke(None, config=config)

        # 执行完成后可以删除工作流状态
        #await workflow.checkpointer.delete_thread(request.thread_id)

        return QueryResponse(
            status="success",
            result=result["result"],
        )
    except Exception as e:
        raise HTTPException(status_code=500, detail=str(e))

五、构建“反思机制”智能体

反思智能体 = 在图里多加一个「反思/审稿」节点:通过就往下走,不通过就带着“批评问题”打回「生成」节点重新生成问题。

它解决什么问题?

  • 普通 Agent:问一次 → 答一次,容易漏点、格式乱、逻辑跳步。
  • 反思模式:生成 → 批评 → 改 → 再批评,用多轮自评提高质量(写作、代码、复杂推理等场景)。

和 ReAct 有啥区别?

        ReAct = Reasoning(推理)+ Acting(行动)

        方式:[思考] → 调工具 → [看结果] → [思考] → 调工具 → … → 最终回答

        反思智能体 = 生成答案 → 批评自己的答案 → 不满意就重写

        方式:任务 → [生成草稿/回答] → [反思:哪里不好] → 通过?结束 : 带着批评再生成

流程长什么样?

  1. 生成节点:根据任务 + 上一轮批评,产出/修改草稿(回答)
  2. 反思节点:只负责挑毛病、给建议,不直接替用户定稿
  3. 条件边:通过继续;未通过打回,并限制最大轮数防死循环

State 里通常放什么?

  • user_request:原始任务
  • draft:当前草稿(节点输出的回答)
  • critique:反思意见(打回时给生成节点用)
  • iteration:已迭代次数(配合上限)

什么时候用 / 不用?

  • 适合:质量比速度重要;有明确评判标准(格式、单测、rubric)
  • 不太适合:简单问答、强实时对话(成本高、延迟大);纯检索更适合 RAG + 工具
Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐