是什么
LangGraph 是用状态图(State Graph)来描述多步智能体流程的编排框架。 用它的效果是:复杂的决策、分支、循环不再藏在 if-else 里,而是变成一张可视化、可单步调试的图。
怎么用
- 先把业务流程拆成节点(一个节点 = 一个明确职责),让每一步的输入输出都有契约。
- 用条件边(Conditional Edge)描述分支与循环,让回退与重试逻辑显式而非散落。
- 把全局状态(State)建模成 TypedDict 或 Pydantic 模型,让每一次状态更新都可追踪。
- 在关键节点接入人工审核断点(Human-in-the-loop),让高风险动作有人类签字才放行。
- 用图遍历日志做事后回放,让失败案例能复现、能归因、能写进回归测试。
架构图
flowchart LR
起点 --> 决策节点
决策节点 --> 工具调用
决策节点 --> 人工审核
工具调用 --> 状态更新
人工审核 --> 状态更新
状态更新 --> 终点
LangGraph Patterns
LangGraph builds stateful multi-step LLM workflows as directed graphs. Each node is a Python function; edges define routing between them.
When to Activate
- Building a multi-step LLM pipeline (research → draft → review → publish)
- Implementing human-in-the-loop interrupts or approval steps
- Designing conditional routing based on LLM output
- Adding persistence/memory to an agent across sessions
- Streaming intermediate results to the client
- Coordinating multiple agents as subgraphs
- Debugging
InvalidUpdateError, cycle errors, or state shape issues
Core Concepts
StateGraph
├── State — TypedDict that flows through every node
├── Nodes — functions: State → State update (partial dict)
├── Edges — unconditional routing A → B
├── Conditional — function decides which node to go to next
└── Checkpointer — persists state between invocations (memory)
Minimal Example
from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages
from langchain_openai import ChatOpenAI
# 1. Define state — Annotated[list, add_messages] appends instead of replacing
class State(TypedDict):
messages: Annotated[list, add_messages]
llm = ChatOpenAI(model="gpt-4o-mini")
# 2. Define a node — receives full state, returns partial update
def chatbot(state: State) -> dict:
return {"messages": [llm.invoke(state["messages"])]}
# 3. Build the graph
graph = (
StateGraph(State)
.add_node("chatbot", chatbot)
.add_edge(START, "chatbot")
.add_edge("chatbot", END)
.compile()
)
# 4. Invoke
result = graph.invoke({"messages": [{"role": "user", "content": "Hello!"}]})
print(result["messages"][-1].content)
State Design
from typing import TypedDict, Annotated
from operator import add
# Annotated reducers control how values merge on update
class ResearchState(TypedDict):
# add_messages: appends new messages, deduplicates by ID
messages: Annotated[list, add_messages]
# add (operator.add): appends items from each node update
sources: Annotated[list[str], add]
# Last-write-wins (default — no annotation needed)
query: str
status: str
final_report: str | None
# Optional fields
error: str | None
Key rule: Nodes return a partial dict — only include keys you want to update. LangGraph merges with the existing state using the reducer.
# Node returns partial — only updates 'status' and 'sources'
def fetch_sources(state: ResearchState) -> dict:
sources = search_web(state["query"])
return {
"sources": sources, # add reducer: appends
"status": "sources_ready",
}
Conditional Routing
from langgraph.graph import StateGraph, START, END
def route_after_llm(state: State) -> str:
"""Return the name of the next node (or END)."""
last_message = state["messages"][-1]
# If the LLM called a tool, go to tools node
if last_message.tool_calls:
return "tools"
# Otherwise finish
return END
graph = StateGraph(State)
graph.add_node("llm", call_llm)
graph.add_node("tools", run_tools)
graph.add_edge(START, "llm")
graph.add_conditional_edges(
"llm", # source node
route_after_llm, # routing function → returns node name
{ # optional: map return values to node names
"tools": "tools",
END: END,
},
)
graph.add_edge("tools", "llm") # loop back after tools
Multiple possible routes
def classify_query(state: State) -> str:
query = state["query"].lower()
if "code" in query: return "code_agent"
if "math" in query: return "math_agent"
return "general_agent"
graph.add_conditional_edges(
"classifier",
classify_query,
["code_agent", "math_agent", "general_agent"], # all possible targets
)
Tool Calling
from langchain_core.tools import tool
from langgraph.prebuilt import ToolNode
@tool
def search_web(query: str) -> str:
"""Search the web for current information."""
return web_search_api(query)
@tool
def calculator(expression: str) -> float:
"""Evaluate a mathematical expression."""
return eval(expression) # use safer eval in production
tools = [search_web, calculator]
tool_node = ToolNode(tools) # pre-built node that runs tools
llm_with_tools = ChatOpenAI(model="gpt-4o-mini").bind_tools(tools)
def call_llm(state: State) -> dict:
return {"messages": [llm_with_tools.invoke(state["messages"])]}
def should_continue(state: State) -> str:
return "tools" if state["messages"][-1].tool_calls else END
graph = StateGraph(State)
graph.add_node("llm", call_llm)
graph.add_node("tools", tool_node)
graph.add_edge(START, "llm")
graph.add_conditional_edges("llm", should_continue)
graph.add_edge("tools", "llm")
Persistence (Checkpointers)
Checkpointers save state after every node so the graph can be paused, resumed, or continued in a new session.
from langgraph.checkpoint.memory import MemorySaver # in-process (dev/test)
from langgraph.checkpoint.postgres import PostgresSaver # production
# In-memory checkpointer
memory = MemorySaver()
graph = StateGraph(State).compile(checkpointer=memory)
# PostgreSQL checkpointer
import psycopg
conn = psycopg.connect("postgresql://user:pass@localhost/db")
checkpointer = PostgresSaver(conn)
graph = StateGraph(State).compile(checkpointer=checkpointer)
# thread_id groups messages into a "conversation" — same ID = same history
config = {"configurable": {"thread_id": "user-123-session-1"}}
# First call — creates new thread
result = graph.invoke({"messages": [HumanMessage("Hello")]}, config=config)
# Second call — continues the same thread
result = graph.invoke({"messages": [HumanMessage("Follow up")]}, config=config)
# Get current state of a thread
snapshot = graph.get_state(config)
print(snapshot.values) # current state
print(snapshot.next) # next node to run (empty if finished)
Human-in-the-Loop (Interrupts)
from langgraph.types import interrupt, Command
# interrupt() pauses the graph and surfaces a value to the caller
def approval_step(state: State) -> dict:
# This raises an interrupt — graph pauses here
human_response = interrupt({
"question": "Should I proceed?",
"context": state["draft"],
})
# Execution resumes here when resumed with a Command
if human_response == "yes":
return {"status": "approved"}
return {"status": "rejected"}
graph = StateGraph(State).compile(
checkpointer=memory,
interrupt_before=["approval_step"], # pause BEFORE this node
# interrupt_after=["draft"], # pause AFTER this node
)
# First invocation — runs until interrupt
result = graph.invoke(input, config=config)
# result contains the interrupt value
# Resume after human provides input
result = graph.invoke(
Command(resume="yes"), # pass human decision
config=config,
)
Streaming
# stream_mode options:
# "values" — full state after each node
# "updates" — partial state update from each node
# "messages"— LLM token-by-token streaming
# Stream full state values
for state in graph.stream(input, config=config, stream_mode="values"):
print(state)
# Stream node updates only
for chunk in graph.stream(input, config=config, stream_mode="updates"):
node_name, update = list(chunk.items())[0]
print(f"Node '{node_name}' updated: {update}")
# Stream LLM tokens (best for chat UI)
async for chunk in graph.astream(input, config=config, stream_mode="messages"):
if hasattr(chunk, "content"):
print(chunk.content, end="", flush=True)
# Async streaming in FastAPI
@router.get("/chat/stream")
async def stream_chat(query: str):
async def generate():
async for chunk in graph.astream(
{"messages": [HumanMessage(query)]},
stream_mode="messages",
):
if hasattr(chunk, "content") and chunk.content:
yield f"data: {chunk.content}\n\n"
return StreamingResponse(generate(), media_type="text/event-stream")
Subgraphs (Multi-Agent)
# Define a specialised sub-agent as its own graph
researcher = (
StateGraph(ResearchState)
.add_node("search", search_web)
.add_node("summarize", summarize)
.add_edge(START, "search")
.add_edge("search", "summarize")
.add_edge("summarize", END)
.compile()
)
writer = (
StateGraph(WriterState)
.add_node("draft", draft_content)
.add_node("refine", refine_draft)
.compile()
)
# Orchestrator graph uses sub-agents as nodes
def run_researcher(state: OrchestratorState) -> dict:
result = researcher.invoke({"query": state["topic"]})
return {"research": result["summary"]}
orchestrator = (
StateGraph(OrchestratorState)
.add_node("research", run_researcher)
.add_node("write", run_writer)
.add_edge(START, "research")
.add_edge("research", "write")
.add_edge("write", END)
.compile(checkpointer=memory)
)
Agentex Integration
In Agentex Temporal agents, LangGraph runs inside a Temporal activity (not directly in the workflow). The ADK provides helpers:
from agentex.lib import adk
# In an activity:
async def run_langgraph_agent(params: AgentParams) -> str:
graph = build_my_graph()
# stream_langgraph_events sends each token/update to the Agentex UI
async for event in adk.stream_langgraph_events(
graph=graph,
inputs={"messages": [HumanMessage(params.user_message)]},
task_id=params.task_id,
):
pass
return final_result
# Checkpointer backed by Agentex state (MongoDB) for persistence
checkpointer = adk.create_checkpointer(task_id=params.task_id)
graph = build_my_graph().compile(checkpointer=checkpointer)
Debugging
# Print the graph structure
print(graph.get_graph().draw_ascii())
# Print state at each step
for step in graph.stream(input, stream_mode="values"):
print("---")
for k, v in step.items():
print(f" {k}: {v}")
# Inspect checkpointed history
history = list(graph.get_state_history(config))
for snapshot in history:
print(snapshot.values, snapshot.next, snapshot.created_at)
# Replay from a specific checkpoint
graph.invoke(None, config={**config, "checkpoint_id": old_checkpoint_id})
Common Errors
| Error | Cause | Fix |
|---|---|---|
InvalidUpdateError |
Node returned a key not in State TypedDict | Add the key to State or remove from return |
GraphRecursionError |
Cycle with no termination condition | Add conditional edge → END when done |
| State not persisting | No checkpointer compiled | Add checkpointer=memory to .compile() |
| Interrupt not working | Missing checkpointer | Interrupts require a checkpointer |
add_messages duplicating |
Returning same message ID twice | Return new messages only; don't re-include history |
Red Flags
- Returning full state from a node — returning the entire state dict instead of a partial update overwrites all fields and breaks reducers; nodes must return only the keys they changed
- Cycles with no exit condition — a loop between two nodes with no conditional edge to
ENDcausesGraphRecursionError; always add a conditional edge that can reachEND MemorySaverin production — in-process memory is lost on worker restart; usePostgresSaver(or another persistent backend) for any deployed graph- Non-deterministic code in node functions — calling
time.time(),random, or direct HTTP requests inside nodes makes replay unpredictable in LangGraph Cloud and Temporal-hosted graphs; use activity patterns for side effects - Missing
thread_idor reusing it across unrelated sessions — reusing athread_idcontinues an old conversation; always generate a unique ID per session and pass it in the config'sconfigurabledict - Human-in-the-loop without a checkpointer —
interrupt()silently does nothing if the graph was compiled without a checkpointer; interrupts requirecheckpointer=in.compile() - Accessing relationship fields across incompatible subgraph state types — parent and subgraph states must have compatible shapes; passing keys the subgraph doesn't declare in its
TypedDictcausesInvalidUpdateError add_messageson a field that isn't a message list — annotating a plain list of strings withadd_messagesdeduplicates by message ID and discards entries without one; useoperator.addfor plain list fields
Checklist
- State is a
TypedDictwith explicit reducers (add_messages,add) for list fields - Nodes return partial dicts — only updated keys
- All cycles have a conditional edge that can route to
END - Tools defined with
@tooldecorator and bound to LLM with.bind_tools() - Production graphs use
PostgresSaver(notMemorySaver) -
thread_idin config is unique per conversation/session - Human-in-the-loop graphs always compiled with a checkpointer
- Streaming uses
astreamfor async contexts