
为什么普通的链式调用不够用用LangChain构建一个简单的RAG问答系统很容易但现实中的AI应用往往更复杂需要根据用户意图走不同的处理路径、需要在某个步骤失败后回退重试、需要让人类在关键节点审批、需要维护跨对话的状态。这时候简单的链式调用A→B→C→D就显得捉襟见肘了。LangGraph用有状态的图结构来解决这个问题——每个节点是一个处理步骤每条边是条件跳转整个工作流的状态在图中持久化和流转。## LangGraph的核心概念### 状态StateState是图中流转的共享数据容器所有节点都读写同一个Statepythonfrom typing import TypedDict, Annotatedfrom langgraph.graph.message import add_messagesclass WorkflowState(TypedDict): # 消息列表使用add_messages reducer自动合并 messages: Annotated[list, add_messages] # 用户意图分类 intent: str | None # 检索到的文档 retrieved_docs: list[dict] # 最终回答 answer: str | None # 错误信息 error: str | None # 重试次数 retry_count: intState的设计是LangGraph工作流设计中最重要的决策。好的State设计原则- 只包含真正需要在节点间共享的数据- 为可能需要累积的字段定义合适的Reducer- 不要把临时计算结果放入State### 节点Node节点是普通的Python函数接收State返回State的更新pythonfrom langchain_openai import ChatOpenAIfrom langchain_core.messages import SystemMessage, HumanMessagellm ChatOpenAI(modelgpt-4o, temperature0)async def classify_intent(state: WorkflowState) - dict: 分类用户意图节点 user_message state[messages][-1].content response await llm.ainvoke([ SystemMessage( 将用户消息分类为以下意图之一 - search: 需要搜索信息 - calculation: 需要进行计算 - chitchat: 日常对话 - unknown: 无法识别 只返回意图标签不要其他内容。 ), HumanMessage(user_message) ]) return {intent: response.content.strip()}async def retrieve_documents(state: WorkflowState) - dict: 文档检索节点 query state[messages][-1].content # 调用向量数据库检索 docs await vector_store.asimilarity_search(query, k5) return { retrieved_docs: [ {content: doc.page_content, metadata: doc.metadata} for doc in docs ] }async def generate_answer(state: WorkflowState) - dict: 生成回答节点 context \n\n.join([ doc[content] for doc in state.get(retrieved_docs, []) ]) response await llm.ainvoke( state[messages] [ SystemMessage(f基于以下上下文回答问题\n{context}) ] ) return { answer: response.content, messages: [response] # 添加到消息历史 }### 边Edge与条件路由边决定了节点执行完后流向哪里pythonfrom langgraph.graph import ENDdef route_by_intent(state: WorkflowState) - str: 根据意图决定下一个节点 intent state.get(intent, unknown) routing { search: retrieve_documents, calculation: calculate, chitchat: generate_chitchat_response, unknown: END # 直接结束 } return routing.get(intent, END)def should_retry(state: WorkflowState) - str: 判断是否需要重试 if state.get(error) and state.get(retry_count, 0) 3: return retrieve_documents # 重试 elif state.get(error): return handle_error # 超过重试次数走错误处理 else: return generate_answer # 正常流程## 构建完整的工作流pythonfrom langgraph.graph import StateGraph, START, ENDfrom langgraph.checkpoint.memory import MemorySaver# 创建图workflow StateGraph(WorkflowState)# 添加节点workflow.add_node(classify_intent, classify_intent)workflow.add_node(retrieve_documents, retrieve_documents)workflow.add_node(generate_answer, generate_answer)workflow.add_node(calculate, calculate_node)workflow.add_node(generate_chitchat_response, chitchat_node)workflow.add_node(handle_error, error_handler_node)# 添加边workflow.add_edge(START, classify_intent) # 起点# 条件边根据意图路由workflow.add_conditional_edges( classify_intent, route_by_intent, { retrieve_documents: retrieve_documents, calculate: calculate, generate_chitchat_response: generate_chitchat_response, END: END })# 检索后判断是否需要重试workflow.add_conditional_edges( retrieve_documents, should_retry, { retrieve_documents: retrieve_documents, # 循环重试 generate_answer: generate_answer, handle_error: handle_error })# 终止节点workflow.add_edge(generate_answer, END)workflow.add_edge(calculate, END)workflow.add_edge(generate_chitchat_response, END)workflow.add_edge(handle_error, END)# 添加检查点用于状态持久化和Human-in-the-loopcheckpointer MemorySaver()# 编译图app workflow.compile(checkpointercheckpointer)## Human-in-the-loop在关键节点插入人工审批这是LangGraph最强大的特性之一。在图中的任何节点可以设置中断点等待人工确认后再继续pythonfrom langgraph.types import interrupt, Commandasync def review_before_publish(state: WorkflowState) - dict: 发布前的人工审批节点 # interrupt()会暂停图的执行等待人工输入 human_decision interrupt({ type: review_request, content: state[answer], question: 这个回答是否可以发送给用户, options: [approve, reject, edit] }) if human_decision[action] approve: return {approved: True} elif human_decision[action] edit: return { answer: human_decision[edited_content], approved: True } else: return {approved: False, error: 被人工拒绝}# 添加审批节点到图workflow.add_node(review_before_publish, review_before_publish)# 编译时指定中断节点app workflow.compile( checkpointercheckpointer, interrupt_before[review_before_publish] # 在这个节点前中断)# 使用示例async def run_with_human_review(): thread_id conversation-123 config {configurable: {thread_id: thread_id}} # 第一次运行会在review节点前暂停 result await app.ainvoke( {messages: [HumanMessage(帮我写一封给客户的邮件)]}, configconfig ) # 此时result包含了待审批的内容 # 用户/管理员查看后做决定 # 恢复执行传入人工决定 final_result await app.ainvoke( Command(resume{action: approve}), configconfig )## 持久化状态跨对话记忆使用PostgreSQL作为检查点存储实现跨会话的状态持久化pythonfrom langgraph.checkpoint.postgres.aio import AsyncPostgresSaverimport psycopgasync def create_app_with_persistence(): conn await psycopg.AsyncConnection.connect( postgresql://user:passlocalhost/langgraph_db ) checkpointer AsyncPostgresSaver(conn) await checkpointer.setup() app workflow.compile(checkpointercheckpointer) return app# 同一个thread_id的多次调用会共享状态# 用户明天继续对话上下文还在async def continue_conversation(user_id: str, message: str): app await create_app_with_persistence() config {configurable: {thread_id: user_id}} result await app.ainvoke( {messages: [HumanMessage(message)]}, configconfig ) return result## 可视化与调试LangGraph提供了内置的可视化工具让你能直观看到工作流结构python# 生成Mermaid图mermaid_code app.get_graph().draw_mermaid()print(mermaid_code)# 或者生成PNG图image_data app.get_graph().draw_mermaid_png()with open(workflow.png, wb) as f: f.write(image_data)调试技巧python# 打印每个节点的输入输出async def debug_run(input_state: dict): async for event in app.astream_events(input_state, versionv2): if event[event] on_chain_start: print(f→ 节点开始: {event[name]}) elif event[event] on_chain_end: print(f✓ 节点结束: {event[name]}) print(f 输出: {event[data][output]})LangGraph的图状态机模型让复杂的AI工作流变得可理解、可维护、可调试。当你的Agent需要处理分支、循环、人机协同时这是目前最成熟的工程方案。