1、概念理解

图的构成:图由节点(node)和边(edge)以及状态(status)组成

1.1、Graph(图)是什么

  • 把一次 AI 应用的执行过程拆成多个“节点(node)”,用“边(edge)”把节点连接成流程。

  • 节点可以是:调用模型、调用工具、做规则判断、读写数据库、检索知识库……

  • 图的好处:流程清晰、可控、可扩展、可观测。

1.2、State(状态)是什么

  • State 是图运行过程中的“共享数据袋子”,所有节点都能读它、并向它写更新。

1.3、Reducer(合并器)是什么:Annotated 的意义就在这

  • 节点函数返回的不是“整个新 state”,而是“对 state 的局部更新(delta)”,例如: {"messages": [ai_message]}

  • LangGraph 要把“旧 state + delta” 合并成“新 state”,就需要合并规则。

  • messages: Annotated[list, add_messages] 的意思就是:

    • messages 这个字段的合并规则用 add_messages

    • 每次节点返回 {"messages": [...]} ,LangGraph 会自动把它“累加”进 state 的 messages 里,而不是覆盖。

完整代码:

import os
from typing import Annotated, TypedDict
from langchain_openai import ChatOpenAI
from langchain_core.messages import HumanMessage, AIMessage
from langgraph.graph import StateGraph, END
from langgraph.graph.message import add_messages
from langgraph.checkpoint.memory import MemorySaver
from dotenv import load_dotenv

load_dotenv()

# 配置 DeepSeek
os.environ["OPENAI_API_KEY"] = os.getenv('DEEPSEEK_API_KEY')
os.environ["OPENAI_API_BASE"] = "https://api.deepseek.com"


class AgentState(TypedDict):
    messages: Annotated[list, add_messages]


def agent_node(state: AgentState):
    """
    支持 token/chunk 流式输出的智能体节点
    """
    model = ChatOpenAI(model="deepseek-chat", temperature=0.7, streaming=True)

    ai_tokens = []
    # 遍历流式返回的 chunk
    for chunk in model.stream(state["messages"]):
        # chunk.content 才是字符串
        token_str = chunk.content
        print(token_str, end="", flush=True)  # 流式打印
        ai_tokens.append(token_str)

    # 拼接成完整消息
    ai_message = AIMessage(content="".join(ai_tokens))
    return {"messages": [ai_message]}


# 创建状态图
workflow = StateGraph(AgentState)
workflow.add_node("agent", agent_node)
workflow.set_entry_point("agent")
workflow.add_edge("agent", END)

memory = MemorySaver()
app = workflow.compile(checkpointer=memory)


def run_stream(user_input: str, config: dict):
    """
    使用 app.stream() 执行一轮对话
    """
    for event in app.stream({"messages": [HumanMessage(content=user_input)]}, config=config):
        if "agent" in event:
            ai_message = event["agent"]["messages"][-1]
            # 只输出 content
            print(f"\n[AI 完整消息] {ai_message.content}")


if __name__ == "__main__":
    thread_id = "会话用户_001"
    config = {"configurable": {"thread_id": thread_id}}

    print("--- Round 1 ---")
    run_stream("我叫小明,我是一名程序员。", config)

    print("\n--- Round 2 ---")
    run_stream("我的职业是什么?", config)

效果展示:

2、环境准备(以 DeepSeek + ChatOpenAI 为例)

在上述代码中使用了这种方式配置:

  • .env 里放 DEEPSEEK_API_KEY

  • 用 dotenv 加载环境变量

  • 把 DeepSeek 当作 OpenAI 兼容接口来用:

    • OPENAI_API_KEY = DEEPSEEK_API_KEY

    • OPENAI_API_BASE = https://api.deepseek.com

在代码里关键行:

  • load_dotenv() :加载 .env

  • os.environ["OPENAI_API_KEY"] = os.getenv('DEEPSEEK_API_KEY')

  • os.environ["OPENAI_API_BASE"] = "https://api.deepseek.com"

常见踩坑:

  • .env 不在当前工作目录时加载不到

  • key 为空时模型会报鉴权错误(可先 print(os.getenv("DEEPSEEK_API_KEY")) 自查)

3、定义“消息结构”和“状态结构”(State)

3.1 消息类型(messages 里放什么)

在上述代码中我们用的是 LangChain 的消息类型:

  • HumanMessage :用户消息

  • AIMessage :模型回复

3.2 定义 State(最关键的“可记忆”写法)

class AgentState(TypedDict):
    messages: Annotated[list, add_messages]

这句带来的行为是:

  • 任何节点只要返回 {"messages": [xxx]}

  • LangGraph 就会用 add_messages 把这部分消息合并进历史 messages 里

  • 你不用自己写 append ,合并发生在“节点执行完、写回状态”的阶段

如果你不写 Annotated[..., add_messages] :

  • 多数情况下 messages 会被“覆盖”,历史对话就不稳定了

4、定义节点(Node)

4.1 节点函数输入输出规则(最重要的接口约定)

一个节点函数一般长这样:

  • 输入: state (当前状态)

  • 输出: dict (要写回 state 的增量更新)

上述代码中节点是 agent_node :

  • 读: state["messages"]

  • 调模型

  • 写回: return {"messages": [ai_message]}

关键点:

  • 返回的是增量更新,不是全量 state

  • messages 的累加靠 reducer( add_messages )完成

4.2 流式输出(streaming)

我们在节点使用stream()来流式输出,得到ai的回复,最后返回AIMessage消息。

def agent_node(state: AgentState):
    """
    支持 token/chunk 流式输出的智能体节点
    """
    model = ChatOpenAI(model="deepseek-chat", temperature=0.7, streaming=True)

    ai_tokens = []
    # 遍历流式返回的 chunk
    for chunk in model.stream(state["messages"]):
        # chunk.content 才是字符串
        token_str = chunk.content
        print(token_str, end="", flush=True)  # 流式打印
        ai_tokens.append(token_str)

    # 拼接成完整消息
    ai_message = AIMessage(content="".join(ai_tokens))
    return {"messages": [ai_message]}

这会让你在终端看到模型一个 token/片段一个片段吐出来。 最后你自己拼成一个完整的 AIMessage 再写回 state:

ai_message = AIMessage(content="".join(ai_tokens))
return {"messages": [ai_message]}

注意:

  • 流式打印只是“展示效果”

  • 真正进入记忆的,是你最终写回 state 的那条 AIMessage

5、搭建图(StateGraph + 节点 + 边)

在我们定义好所有的节点之后,我们就需要来进行图的构建,本章只展示了单节点的构建,对于复杂的节点需要使用到条件边和对应路由函数 :

workflow = StateGraph(AgentState)
workflow.add_node("agent", agent_node)
workflow.set_entry_point("agent")
workflow.add_edge("agent", END)

解释一下每句在干嘛:

  • StateGraph(AgentState) :声明这张图运行时 state 长什么样

  • add_node("agent", agent_node) :注册一个叫 agent 的节点

  • set_entry_point("agent") :图从哪个节点开始跑

  • add_edge("agent", END) : agent 跑完就结束

入门阶段你可以记住一个最小图套路:

  • 一个节点:模型回答

  • 一条边:节点 → END

6、选择记忆(Checkpointer)并编译图(compile)

6.1 为什么要 compile

  • compile 会把“图的结构 + 合并规则 + checkpointer 等运行配置”变成一个可执行的 app

  • 之后你就用 app.invoke() / app.stream() 跑图

6.2 MemorySaver:最简单的记忆(仅进程内)

我们可以使用MemorySaver来进行简单的记忆 :

memory = MemorySaver()
app = workflow.compile
(checkpointer=memory)

它会把每个 thread_id 的 state 存在内存里:

  • 程序不退出:记忆在

  • 程序一退出:全没了(这是正常的)

7、运行图(invoke / stream)+ thread_id(让记忆“串起来”)

7.1 thread_id 的意义:同一个 thread_id 才能续上历史

我们需要填入thread_id,每一个thread_id代表一个用户,拥有它自己的记忆。如果我们不填写的话,则每次都会随机生成,导致每次会话都是新的记忆空间。

thread_id = "会话用户_001"
config = {"configurable": 
{"thread_id": thread_id}}

要点:

  • thread_id 是“记忆的索引键”

  • 同一个 thread_id:会加载上一次保存的 state,然后继续合并

  • 换一个 thread_id:就相当于新开一段对话

7.2 运行方式:app.stream(带事件流)

在上述代码中我们封装了 run_stream :

for event in app.stream({"messages": [HumanMessage(content=user_input)]}, config=config):
    if "agent" in event:
        ai_message = event["agent"]
        ["messages"][-1]
        print(f"\n[AI 完整消息] 
        {ai_message.content}")

解释一下这里发生了什么:

  • 传入 {"messages": [HumanMessage(...)]} 作为“本轮输入增量”

  • LangGraph 会:

    • 用 thread_id 取出旧 state(含历史 messages)

    • 用 add_messages 把这条 HumanMessage 合并进去

    • 运行 agent_node

    • agent_node 返回 {"messages": [AIMessage(...)]}

    • 再次合并进 state

    • 保存 state 到 checkpointer(MemorySaver)

所以“自动累加”的关键点是: 每一步的输出都是 delta,合并由 reducer(Annotated 指定)统一完成 。

Logo

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

更多推荐