目 录CONTENT

文章目录

LangGraph 流式输出与调试:Streaming、Debug 与 LangGraph Studio

PySuper
2025-10-11 / 0 评论 / 0 点赞 / 1 阅读 / 0 字
温馨提示:
本文最后更新于2026-05-22,若内容或图片失效,请留言反馈。 所有牛逼的人都有一段苦逼的岁月。 但是你只要像SB一样去坚持,终将牛逼!!! ✊✊✊

作者:PySuper | 来源:zhengxingtao.com | 更新日期:2026-03-22

前面我们从入门到生产部署,把 LangGraph 的核心能力都过了一遍。但有个问题一直没深聊——你的 Agent 跑起来之后,用户在那干等,啥也看不到

invoke() 一跑就是 30 秒,用户以为系统卡死了,直接关页面走人。这不行。

本篇聚焦 LangGraph 的实时交互能力调试能力——这是从"能跑"到"好用"的关键跨越。

一、为什么需要流式输出

1.1 invoke 的致命问题

python

# 用户提交问题后,只能干等
result = graph.invoke({"messages": [("user", "帮我分析这份财报")]})
# ... 30 秒后才有结果
# 用户内心:是不是挂了?

invoke() 的核心问题:全程黑盒。工作流内部可能经历"检索文档→分析数据→生成报告"等多个步骤,但用户只看到请求发出,然后等一个漫长的最终结果。

1.2 流式的价值

plaintext

┌──────────────────────────────────────────────────────┐
│                    用户体验对比                        │
├──────────────────────────────────────────────────────┤
│                                                      │
│  invoke 模式:                                        │
│  ┌─────┐                              ┌──────────┐  │
│  │请求  │ ████████████████████████████ │ 最终结果  │  │
│  └─────┘        30s 黑盒等待           └──────────┘  │
│                                                      │
│  stream 模式:                                        │
│  ┌─────┐ ┌──────┐ ┌──────┐ ┌──────┐ ┌──────────┐   │
│  │请求  │ │检索中 │ │分析中 │ │撰写中 │ │ 完成 ✓   │   │
│  └─────┘ └──────┘ └──────┘ └──────┘ └──────────┘   │
│           ↑ 实时反馈,用户知道系统在干活               │
│                                                      │
└──────────────────────────────────────────────────────┘

流式输出的三大价值:

  • 实时反馈:用户知道系统在工作,不是卡了

  • 进度可见:能看到当前执行到哪一步了

  • 体验提升:逐 Token 输出(打字机效果),体感延迟大幅降低

1.3 不同场景对流式的需求差异

表格

场景

核心需求

推荐模式

ChatBot 对话

逐字输出,打字机效果

messages

工作流监控

步骤进度、当前节点

updates

状态审计

每步完整状态快照

values

长任务进度

百分比、中间结果

custom

开发调试

底层运行时事件

debug

二、四种流式模式详解

LangGraph 的 stream()astream() 方法支持通过 stream_mode 参数选择不同的流式模式,还可以组合使用。

plaintext

┌──────────────────────────────────────────────────────────────┐
│                   LangGraph 流式模式全景                       │
├──────────────┬───────────────────┬───────────────────────────┤
│    模式       │    返回内容        │      典型场景              │
├──────────────┼───────────────────┼───────────────────────────┤
│  values      │ 每步完整 State     │ 全量状态监控、审计         │
│  updates     │ 每步增量更新       │ 步骤进度展示              │
│  messages    │ 逐 Token 输出     │ ChatBot 打字机效果        │
│  custom      │ 自定义数据推送     │ 进度百分比、业务日志       │
│  debug       │ 底层运行时事件     │ 开发调试、问题排查         │
└──────────────┴───────────────────┴───────────────────────────┘

先准备一个通用的示例图,后续所有模式都基于它来演示:

python

from typing import TypedDict, Annotated
from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages
from langchain_core.messages import HumanMessage, AIMessage
from langchain_openai import ChatOpenAI

llm = ChatOpenAI(model="gpt-4o-mini")


class State(TypedDict):
    messages: Annotated[list, add_messages]
    topic: str
    joke: str


def refine_topic(state: State) -> dict:
    """优化用户输入的主题"""
    response = llm.invoke(
        f"将以下主题优化为更适合生成笑话的表述,只返回优化后的主题:{state['topic']}"
    )
    return {"topic": response.content}


def generate_joke(state: State) -> dict:
    """根据主题生成笑话"""
    response = llm.invoke(f"讲一个关于{state['topic']}的笑话")
    return {"joke": response.content, "messages": [AIMessage(content=response.content)]}


# 构建图
graph = (
    StateGraph(State)
    .add_node("refine_topic", refine_topic)
    .add_node("generate_joke", generate_joke)
    .add_edge(START, "refine_topic")
    .add_edge("refine_topic", "generate_joke")
    .add_edge("generate_joke", END)
    .compile()
)

2.1 values 模式

每步返回完整 State——相当于给每个 superstep 拍了张全量快照。

python

# values 模式:每步之后输出完整状态
for chunk in graph.stream(
    {"topic": "程序员"},
    stream_mode="values",
):
    print(chunk)
    print("---")

输出:

plaintext

{'topic': '程序员', 'messages': [], 'joke': ''}
---
{'topic': '程序员的日常生活', 'messages': [], 'joke': ''}
---
{'topic': '程序员的日常生活', 'messages': [AIMessage(...)], 'joke': '为什么程序员...'}
---

可以看到,每一步都会返回整个 State 的当前值。即使某一步只改了 topic,你拿到的也是包含所有字段的完整状态。

适用场景:全量状态监控、状态审计、需要完整上下文的前端渲染。

⚠️ 注意:如果 State 很大(比如长对话历史),values 模式每步都传全量数据,网络开销不可忽视。对生产环境推荐用 updates 模式。

2.2 updates 模式

只返回增量更新——每个节点只返回它自己修改的那部分字段。

python

# updates 模式:只看增量
for chunk in graph.stream(
    {"topic": "程序员"},
    stream_mode="updates",
):
    for node, update in chunk.items():
        print(f"节点 [{node}] 更新: {update}")

输出:

plaintext

节点 [refine_topic] 更新: {'topic': '程序员的日常生活'}
节点 [generate_joke] 更新: {'joke': '为什么程序员...', 'messages': [AIMessage(...)]}

对比 values 模式:

plaintext

┌────────────────────────────────────────────────────┐
│              values vs updates 对比                  │
├────────────────────────────────────────────────────┤
│                                                    │
│  values:                                           │
│  Step 0 → {topic: "程序员", joke: "", msgs: []}    │
│  Step 1 → {topic: "XX",    joke: "", msgs: []}    │
│  Step 2 → {topic: "XX",    joke: "YY", msgs: [Z]} │
│           ↑ 每次都是完整状态                          │
│                                                    │
│  updates:                                          │
│  Step 1 → {topic: "XX"}       ← 只改了 topic      │
│  Step 2 → {joke: "YY", msgs: [Z]} ← 只改了这俩    │
│           ↑ 每次只有增量                              │
│                                                    │
└────────────────────────────────────────────────────┘

适用场景:步骤进度展示、前端只需知道"哪一步做了什么"。

2.3 messages 模式

逐 Token 输出——这是前端最常用的模式,实现"打字机效果"。

messages 模式返回的每个 chunk 是一个元组:(message_chunk, metadata),其中 message_chunkAIMessageChunk 对象,metadata 包含当前节点等信息。

python

# messages 模式:逐 Token 输出
for chunk, metadata in graph.stream(
    {"topic": "程序员"},
    stream_mode="messages",
):
    # 只输出模型节点的 token
    if metadata.get("langgraph_node") == "generate_joke":
        print(chunk.content, end="", flush=True)

输出效果:

plaintext

为什么程序员总是分不清万圣节和圣诞节?因为 Oct 31 = Dec 25。

一个字一个字蹦出来——这就是 ChatGPT 那种打字机效果。

前端 SSE 集成

messages 模式最常和前端 SSE(Server-Sent Events)搭配使用。基本思路:

plaintext

┌────────────┐    SSE 流     ┌─────────────┐
│  LangGraph  │ ──────────→  │   FastAPI    │
│   astream   │  token流     │  SSE 推送    │
└────────────┘               └──────┬──────┘
                                     │ SSE
                              ┌──────▼──────┐
                              │   前端       │
                              │ EventSource │
                              └─────────────┘

前端用 EventSource 接收:

javascript

const eventSource = new EventSource("/api/chat/stream?message=你好");
eventSource.onmessage = (event) => {
    const data = JSON.parse(event.data);
    if (data.type === "token") {
        appendToChat(data.content);  // 逐字追加到聊天界面
    } else if (data.type === "done") {
        eventSource.close();
    }
};

💡 提示messages 模式要求 State 中必须有 messages 字段,且使用 add_messages reducer。如果你的 State 没有消息列表,用这个模式不会得到任何输出。

2.4 custom 模式

节点内推送自定义数据——这是最灵活的模式,可以在节点执行过程中推送任意数据。

使用 get_stream_writer() 获取写入器,在节点内主动推送自定义数据:

python

import time
from langgraph.config import get_stream_writer


class FileProcessState(TypedDict):
    files: list[str]
    results: list[str]


def process_files(state: FileProcessState) -> dict:
    """批量处理文件,通过 custom 模式推送进度"""
    writer = get_stream_writer()
    files = state["files"]
    results = []

    for i, file in enumerate(files, 1):
        # 推送自定义进度信息
        writer({"progress": f"正在处理第 {i}/{len(files)} 个文件: {file}"})
        writer({"percentage": round(i / len(files) * 100, 1)})

        # 模拟文件处理
        time.sleep(0.5)
        results.append(f"{file}: 处理完成")

    writer({"progress": "所有文件处理完毕 ✓"})

    return {"results": results}


# 构建图
process_graph = (
    StateGraph(FileProcessState)
    .add_node("process", process_files)
    .add_edge(START, "process")
    .add_edge("process", END)
    .compile()
)

# 消费 custom 流
for chunk in process_graph.stream(
    {"files": ["report.pdf", "data.csv", "config.yaml", "output.xlsx"]},
    stream_mode="custom",
):
    print(chunk)

输出:

plaintext

{'progress': '正在处理第 1/4 个文件: report.pdf'}
{'percentage': 25.0}
{'progress': '正在处理第 2/4 个文件: data.csv'}
{'percentage': 50.0}
{'progress': '正在处理第 3/4 个文件: config.yaml'}
{'percentage': 75.0}
{'progress': '正在处理第 4/4 个文件: output.xlsx'}
{'percentage': 100.0}
{'progress': '所有文件处理完毕 ✓'}

适用场景

  • 长时间任务的进度百分比

  • 业务日志推送("正在连接数据库...")

  • 中间结果预览

  • 任何不属于 State 但需要实时传递给前端的数据

⚠️ 关键custom 模式需要你在节点代码中主动调用 writer() 推送数据。如果你忘了调用,stream_mode="custom" 不会有任何输出!

多模式组合

stream_mode 支持传入列表,同时消费多种模式:

python

# 同时获取 updates 和 custom
for mode, chunk in graph.stream(
    input_data,
    stream_mode=["updates", "custom"],
):
    if mode == "updates":
        print(f"[状态更新] {chunk}")
    elif mode == "custom":
        print(f"[自定义] {chunk}")

组合模式时,每个 chunk 会以 (mode_name, data) 元组的形式返回。

2.5 debug 模式

底层运行时事件——让你看到 LangGraph 引擎内部的每一次调度、每一条 Channel 写入。

python

# debug 模式:看引擎内部发生了什么
for chunk in graph.stream(
    {"topic": "程序员"},
    stream_mode="debug",
):
    print(chunk)

输出示例(简化):

plaintext

{'type': 'task', 'timestamp': '...', 'step': 1, 'payload': {'id': '...', 'name': 'refine_topic', 'input': {...}, 'triggers': ['__start__']}}
{'type': 'task_result', 'timestamp': '...', 'step': 1, 'payload': {'id': '...', 'name': 'refine_topic', 'output': {'topic': '...'}}}
{'type': 'task', 'timestamp': '...', 'step': 2, 'payload': {'id': '...', 'name': 'generate_joke', 'input': {...}, 'triggers': ['refine_topic']}}
{'type': 'task_result', 'timestamp': '...', 'step': 2, 'payload': {'id': '...', 'name': 'generate_joke', 'output': {'joke': '...'}}}

debug 模式的事件类型:

plaintext

┌──────────────────────────────────────────────────────┐
│               debug 模式事件类型                       │
├────────────────┬─────────────────────────────────────┤
│   事件类型      │          说明                       │
├────────────────┼─────────────────────────────────────┤
│  task          │ 节点被调度执行                       │
│  task_result   │ 节点执行完成,返回结果               │
│  error         │ 节点执行出错                        │
└────────────────┴─────────────────────────────────────┘

适用场景:开发阶段排查问题、理解图的执行流程、验证条件边的路由逻辑。

三、异步流式输出

生产环境几乎都是异步的,LangGraph 提供了 astream()astream_events() 两个异步流式接口。

3.1 astream 基础用法

astream()stream() 的异步版本,用法完全一致:

python

import asyncio


async def stream_demo():
    # 异步流式输出
    async for chunk, metadata in graph.astream(
        {"topic": "程序员"},
        stream_mode="messages",
    ):
        node = metadata.get("langgraph_node", "")
        if chunk.content:
            print(chunk.content, end="", flush=True)


asyncio.run(stream_demo())

3.2 astream_events 更细粒度

astream_events() 提供了更细粒度的事件流,能捕获 LLM 调用的 start/end、工具调用的输入输出等:

python

async def events_demo():
    async for event in graph.astream_events(
        {"topic": "程序员"},
        version="v2",
    ):
        kind = event["event"]
        if kind == "on_chat_model_stream":
            # LLM Token 流
            token = event["data"]["chunk"].content
            if token:
                print(token, end="", flush=True)
        elif kind == "on_chain_start":
            print(f"\n[开始] {event['name']}")
        elif kind == "on_chain_end":
            print(f"\n[结束] {event['name']}")


asyncio.run(events_demo())

⚠️ 注意astream_events 的事件类型非常多(on_chat_model_starton_chat_model_streamon_chat_model_endon_tool_starton_tool_end 等),需要仔细过滤。建议只关注你关心的几种事件类型。

3.3 与 FastAPI 集成的完整示例

这是一个完整可运行的 FastAPI + LangGraph 流式服务:

python

# server.py
import json
import uuid
import asyncio
from typing import TypedDict, Annotated

from fastapi import FastAPI
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import StreamingResponse
from pydantic import BaseModel

from langgraph.graph import StateGraph, START, END
from langgraph.graph.message import add_messages
from langgraph.checkpoint.memory import InMemorySaver
from langchain_core.messages import HumanMessage, AIMessage
from langchain_openai import ChatOpenAI

# ─── 状态定义 ───
class ChatState(TypedDict):
    messages: Annotated[list, add_messages]


# ─── 节点定义 ───
llm = ChatOpenAI(model="gpt-4o-mini")


def chatbot(state: ChatState) -> dict:
    response = llm.invoke(state["messages"])
    return {"messages": [response]}


# ─── 构建图 ───
checkpointer = InMemorySaver()

graph = (
    StateGraph(ChatState)
    .add_node("chatbot", chatbot)
    .add_edge(START, "chatbot")
    .add_edge("chatbot", END)
    .compile(checkpointer=checkpointer)
)

# ─── FastAPI 应用 ───
app = FastAPI(title="LangGraph Streaming Chat")

app.add_middleware(
    CORSMiddleware,
    allow_origins=["*"],
    allow_methods=["*"],
    allow_headers=["*"],
)


class ChatRequest(BaseModel):
    message: str
    thread_id: str = "default"


@app.post("/api/chat/stream")
async def chat_stream(req: ChatRequest):
    """SSE 流式聊天接口"""

    async def event_generator():
        config = {"configurable": {"thread_id": req.thread_id}}
        input_data = {"messages": [HumanMessage(content=req.message)]}

        async for chunk, metadata in graph.astream(
            input_data,
            config=config,
            stream_mode="messages",
        ):
            # 只转发模型节点的 token
            if metadata.get("langgraph_node") != "chatbot":
                continue

            content = getattr(chunk, "content", None)
            if not content:
                continue

            # SSE 格式推送
            data = json.dumps({"type": "token", "content": content})
            yield f"data: {data}\n\n"

        # 发送完成信号
        yield f"data: {json.dumps({'type': 'done'})}\n\n"

    return StreamingResponse(
        event_generator(),
        media_type="text/event-stream",
        headers={
            "Cache-Control": "no-cache",
            "Connection": "keep-alive",
            "X-Accel-Buffering": "no",  # 关闭 Nginx 缓冲
        },
    )


@app.post("/api/chat/invoke")
async def chat_invoke(req: ChatRequest):
    """非流式接口(对比用)"""
    config = {"configurable": {"thread_id": req.thread_id}}
    input_data = {"messages": [HumanMessage(content=req.message)]}

    result = await graph.ainvoke(input_data, config=config)
    last_msg = result["messages"][-1]
    return {"response": last_msg.content}

启动服务:

bash

pip install fastapi uvicorn langgraph langchain-openai
uvicorn server:app --reload --port 8000

前端接收:

javascript

// 前端 EventSource 接收 SSE 流
function sendMessage(message) {
    const eventSource = new EventSource(
        `/api/chat/stream?message=${encodeURIComponent(message)}&thread_id=thread-1`
    );

    // 注意:POST 请求不能直接用 EventSource
    // 这里用 fetch + ReadableStream 替代
    fetch("/api/chat/stream", {
        method: "POST",
        headers: { "Content-Type": "application/json" },
        body: JSON.stringify({ message, thread_id: "thread-1" }),
    }).then(async (response) => {
        const reader = response.body.getReader();
        const decoder = new TextDecoder();

        while (true) {
            const { done, value } = await reader.read();
            if (done) break;

            const text = decoder.decode(value);
            const lines = text.split("\n");

            for (const line of lines) {
                if (line.startsWith("data: ")) {
                    const data = JSON.parse(line.slice(6));
                    if (data.type === "token") {
                        appendToChat(data.content);  // 逐字追加
                    } else if (data.type === "done") {
                        console.log("流式输出完成");
                    }
                }
            }
        }
    });
}

3.4 同时推送多种模式

生产环境常见需求:既要 Token 流,又要进度/状态更新。可以用多模式组合:

python

@app.post("/api/chat/stream-full")
async def chat_stream_full(req: ChatRequest):
    """完整流式:Token + 状态更新"""

    async def event_generator():
        config = {"configurable": {"thread_id": req.thread_id}}
        input_data = {"messages": [HumanMessage(content=req.message)]}

        async for mode, chunk in graph.astream(
            input_data,
            config=config,
            stream_mode=["messages", "updates"],
        ):
            if mode == "messages":
                chunk_content, metadata = chunk
                if metadata.get("langgraph_node") == "chatbot" and chunk_content.content:
                    data = json.dumps({"type": "token", "content": chunk_content.content})
                    yield f"data: {data}\n\n"

            elif mode == "updates":
                for node, update in chunk.items():
                    data = json.dumps({"type": "step", "node": node})
                    yield f"data: {data}\n\n"

        yield f"data: {json.dumps({'type': 'done'})}\n\n"

    return StreamingResponse(event_generator(), media_type="text/event-stream")

四、流式与 Checkpoint 的配合

4.1 流式输出时 Checkpoint 自动保存

当你同时使用了 Checkpointer 和流式输出,Checkpoint 的保存是完全自动的——每完成一个 superstep,状态就会被持久化,流式输出和持久化互不干扰:

python

from langgraph.checkpoint.memory import InMemorySaver

checkpointer = InMemorySaver()
graph_with_cp = graph.compile(checkpointer=checkpointer)

# 流式输出的同时,每个 step 的 checkpoint 都自动保存了
config = {"configurable": {"thread_id": "thread-1"}}
for chunk in graph_with_cp.stream({"topic": "程序员"}, config=config, stream_mode="updates"):
    print(chunk)

# 流式结束后,可以查看完整的状态历史
state = graph_with_cp.get_state(config)
print(f"当前步骤: {state.metadata['step']}")

plaintext

┌──────────────────────────────────────────────────────┐
│           流式输出与 Checkpoint 协作流程               │
├──────────────────────────────────────────────────────┤
│                                                      │
│  Step 1: refine_topic 执行                           │
│    ├── 流式推送: {topic: "XX"}                        │
│    └── Checkpoint 自动保存 ✓                         │
│                                                      │
│  Step 2: generate_joke 执行                          │
│    ├── 流式推送: {joke: "YY"}                        │
│    └── Checkpoint 自动保存 ✓                         │
│                                                      │
│  即使中途崩溃,下次用相同 thread_id 即可恢复            │
│                                                      │
└──────────────────────────────────────────────────────┘

4.2 中断时的流式行为

当图执行到 interrupt_beforeinterrupt_after 指定的节点时,流会正常结束,同时 Checkpoint 会保存中断前的状态:

python

# 编译时指定中断点
graph_with_interrupt = graph.compile(
    checkpointer=checkpointer,
    interrupt_before=["generate_joke"],  # 在生成笑话前暂停
)

config = {"configurable": {"thread_id": "thread-2"}}

# 第一次执行:会在 generate_joke 之前中断
for chunk in graph_with_interrupt.stream(
    {"topic": "程序员"}, config=config, stream_mode="updates"
):
    print(chunk)
# 输出: {'refine_topic': {'topic': '程序员的日常生活'}}
# 流正常结束,但 generate_joke 还没执行

# 查看中断状态
state = graph_with_interrupt.get_state(config)
print(f"下个节点: {state.next}")  # ('generate_joke',)

# 恢复执行
for chunk in graph_with_interrupt.stream(None, config=config, stream_mode="updates"):
    print(chunk)
# 输出: {'generate_joke': {'joke': '...'}}

4.3 长时间运行工作流的流式监控

对于可能运行几分钟甚至几小时的工作流,结合 Checkpointer 和流式输出可以实现断点续传 + 实时监控

python

from langgraph.checkpoint.sqlite import SqliteSaver

# 使用 SQLite 做持久化(生产环境推荐 PostgresSaver)
with SqliteSaver.from_conn_string("workflow.db") as saver:
    long_graph = long_workflow.compile(checkpointer=saver)

    config = {"configurable": {"thread_id": "long-job-001"}}

    # 即使连接断开,checkpoint 已保存
    # 重连后用相同 thread_id 恢复
    for chunk in long_graph.stream(
        initial_input, config=config, stream_mode=["updates", "custom"]
    ):
        mode, data = chunk
        if mode == "updates":
            print(f"[步骤] {data}")
        elif mode == "custom":
            print(f"[进度] {data.get('percentage', 0)}%")

五、Debug 调试方法

5.1 print 输出调试

最原始但最有效的方式——在节点函数中 print state:

python

def my_node(state: State) -> dict:
    print(f"===== 进入节点 my_node =====")
    print(f"当前 state: {state}")
    print(f"messages 数量: {len(state.get('messages', []))}")

    result = do_something(state)

    print(f"本节点返回: {result}")
    print(f"===== 退出节点 my_node =====")
    return result

优点:零配置,任何环境都能用

缺点:输出混乱,生产环境不能留,信息不好筛选

💡 进阶技巧:可以用 logging 替代 print,这样可以通过日志级别控制开关:

python

import logging
logger = logging.getLogger("langgraph.debug")


def my_node(state: State) -> dict:
    logger.debug("State in my_node: %s", state)
    # ...

5.2 LangGraph Studio

LangGraph Studio 是 LangChain 官方的可视化调试工具,提供图的可视化渲染、逐步执行、检查点回溯、修改 State 重新执行等能力。

安装与启动

bash

# 安装 CLI
pip install langgraph-cli

# 在项目目录下启动开发服务器
langgraph dev
# 默认启动在 http://localhost:2024

然后在浏览器打开 smith.langchain.com,连接本地服务器即可。

核心功能

plaintext

┌─────────────────────────────────────────────────────────────┐
│                  LangGraph Studio 界面                       │
├──────────┬──────────────────────────────┬───────────────────┤
│          │                              │                   │
│  导航栏   │      画布区域               │   右侧面板        │
│          │                              │                   │
│ 项目列表  │  ┌──────┐   ┌──────────┐   │  状态面板         │
│ 会话历史  │  │start │──→│refine_topic│  │  (查看/修改)      │
│ 检查点    │  └──────┘   └─────┬────┘   │                   │
│          │                    │        │  节点详情          │
│          │              ┌─────▼────┐   │  (输入/输出/耗时)  │
│          │              │gen_joke  │   │                   │
│          │              └─────┬────┘   │  控制台            │
│          │                    │        │  (执行日志)        │
│          │              ┌─────▼────┐   │                   │
│          │              │   end    │   │                   │
│          │              └──────────┘   │                   │
│          │                              │                   │
├──────────┴──────────────────────────────┴───────────────────┤
│  ✅ 绿色 = 执行成功   🔴 红色 = 执行异常   🟡 黄色 = 执行中  │
└─────────────────────────────────────────────────────────────┘

四大核心能力

  1. 图的可视化渲染:自动绘制节点和边,条件边标注判断条件

  2. 逐步执行和检查点回溯:可以回退到任意历史检查点,从那一步重新执行

  3. 修改 State 重新执行:直接编辑状态值,验证"如果这个字段是 X,会走哪条分支"

  4. 分支模拟:手动指定条件边的跳转方向,不用改代码

典型调试流程

plaintext

发现问题 → Studio 中定位异常节点 → 检查输入/输出 → 修改 State 重跑 → 确认修复 → 改代码

5.3 LangSmith Trace

LangSmith 是 LangChain 生态的观测平台,与 LangGraph 深度集成,提供完整的执行追踪能力。

配置集成

python

# 方式一:环境变量(推荐)
import os
os.environ["LANGSMITH_TRACING"] = "true"
os.environ["LANGSMITH_ENDPOINT"] = "https://api.smith.langchain.com"
os.environ["LANGSMITH_API_KEY"] = "ls_xxxxxx"
os.environ["LANGSMITH_PROJECT"] = "my-langgraph-project"

# 方式二:代码中配置
from langchain.callbacks import LangSmithCallbackHandler

handler = LangSmithCallbackHandler(
    project_name="my-langgraph-project",
    api_key="ls_xxxxxx",
)

result = graph.invoke(input_data, config={"callbacks": [handler]})

Trace 提供的信息

plaintext

┌─────────────────────────────────────────────────────────────┐
│                LangSmith Trace 信息                          │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  每个 Node 的输入/输出                                       │
│  ├── refine_topic: input={topic: "程序员"}                  │
│  │                 output={topic: "程序员的日常生活"}         │
│  └── generate_joke: input={topic: "程序员的日常生活"}        │
│                     output={joke: "..."}                     │
│                                                             │
│  执行耗时                                                    │
│  ├── refine_topic: 1.2s                                     │
│  └── generate_joke: 2.8s                                    │
│                                                             │
│  Token 消耗统计                                              │
│  ├── refine_topic: input=45 tokens, output=12 tokens        │
│  └── generate_joke: input=38 tokens, output=67 tokens       │
│  Total: 162 tokens                                          │
│                                                             │
│  错误追踪                                                    │
│  └── 如果某节点报错,完整的异常堆栈                           │
│                                                             │
└─────────────────────────────────────────────────────────────┘

LangSmith 的核心价值:不需要改业务代码,只需要配置环境变量,就能自动采集全链路的执行数据。生产环境必备。

5.4 自定义回调

如果你想在节点执行前后插入自定义逻辑(日志记录、指标采集、告警),可以实现自定义回调:

python

import time
import logging
from langchain_core.callbacks import BaseCallbackHandler

logger = logging.getLogger("langgraph.custom")


class GraphMetricsCallback(BaseCallbackHandler):
    """自定义回调:采集节点执行耗时和 Token 消耗"""

    def __init__(self):
        self._start_times = {}
        self._token_counts = {"input": 0, "output": 0}

    def on_llm_start(self, serialized, prompts, **kwargs):
        """LLM 调用开始"""
        run_id = kwargs.get("run_id")
        self._start_times[run_id] = time.time()
        logger.info("LLM 调用开始: model=%s", serialized.get("name", "unknown"))

    def on_llm_end(self, response, **kwargs):
        """LLM 调用结束"""
        run_id = kwargs.get("run_id")
        if run_id in self._start_times:
            elapsed = time.time() - self._start_times.pop(run_id)
            logger.info("LLM 调用完成: 耗时=%.2fs", elapsed)

        # 统计 Token
        if hasattr(response, "llm_output") and response.llm_output:
            usage = response.llm_output.get("token_usage", {})
            self._token_counts["input"] += usage.get("prompt_tokens", 0)
            self._token_counts["output"] += usage.get("completion_tokens", 0)

    def on_llm_error(self, error, **kwargs):
        """LLM 调用出错"""
        logger.error("LLM 调用失败: %s", str(error))

    def get_metrics(self) -> dict:
        return {
            "total_tokens": self._token_counts["input"] + self._token_counts["output"],
            "input_tokens": self._token_counts["input"],
            "output_tokens": self._token_counts["output"],
        }


# 使用回调
callback = GraphMetricsCallback()
result = graph.invoke(
    {"topic": "程序员"},
    config={"callbacks": [callback]},
)
print(f"本次执行指标: {callback.get_metrics()}")

六、常见调试场景

6.1 Agent 死循环

现象 :Agent 一直在循环调用工具,永远停不下来。

调试步骤

python

# Step 1: 用 debug 模式看执行流程
for chunk in graph.stream(input_data, stream_mode="debug"):
    if chunk["type"] == "task":
        print(f"[Step {chunk['step']}] 调度: {chunk['payload']['name']}")

如果你看到同一个节点反复被调度,就是死循环了。

python

# Step 2: 用 updates 模式看每步的状态变化
for chunk in graph.stream(input_data, stream_mode="updates"):
    for node, update in chunk.items():
        print(f"[{node}] {update}")

常见原因与修复

表格

原因

排查方法

修复方式

条件边的终止条件写错

检查路由函数的返回值

修正 should_continue 的逻辑

LLM 反复调用同一个工具

updates 中 tool_calls

添加最大迭代次数限制

State 中的标志位没更新

values 模式

确认节点正确更新了循环控制字段

python

# 修复:添加最大迭代次数
from langgraph.graph import StateGraph, START, END

def should_continue(state: State) -> str:
    # 安全阀:最多循环 5 次
    if state.get("iteration_count", 0) >= 5:
        return END
    if state.get("task_complete"):
        return END
    return "agent"

6.2 工具调用失败

现象 :工具执行出错,但错误信息被吞了。

调试步骤

python

# Step 1: LangSmith 中查看 tool 节点的完整输入输出
# 如果没有 LangSmith,用 debug 模式
for chunk in graph.stream(input_data, stream_mode="debug"):
    if chunk["type"] == "task_result":
        name = chunk["payload"]["name"]
        output = chunk["payload"]["output"]
        print(f"[{name}] output: {output}")
        if "error" in str(output).lower():
            print(f"  ⚠️ 发现错误!")

python

# Step 2: 在工具节点中加异常捕获
def my_tool(state: State) -> dict:
    try:
        result = call_external_api(state["query"])
        return {"result": result}
    except Exception as e:
        # 把错误信息写进 State,而不是让它被吞掉
        return {
            "result": None,
            "error": f"工具调用失败: {type(e).__name__}: {str(e)}",
        }

6.3 State 意外被覆盖

现象 :某个节点的输出把别的节点设置的字段给覆盖了。

调试步骤

python

# Step 1: 用 values 模式对比每步的完整 State
prev_state = None
for chunk in graph.stream(input_data, stream_mode="values"):
    if prev_state is not None:
        # 对比差异
        for key in set(list(prev_state.keys()) + list(chunk.keys())):
            old_val = prev_state.get(key)
            new_val = chunk.get(key)
            if old_val != new_val:
                print(f"[变化] {key}: {old_val} → {new_val}")
    prev_state = chunk

常见原因

  1. 节点返回了不必要的字段 :节点返回了 {"topic": "new", "joke": ""},空字符串把上一步的结果清零了

  2. Reducer 用错 messages 字段用了 add_messages 但别的字段没加 reducer,默认是覆盖

python

# 修复:确保节点只返回自己需要修改的字段
def refine_topic(state: State) -> dict:
    # ✅ 只返回 topic 的更新
    return {"topic": "优化后的主题"}

# ❌ 不要返回空的其他字段
def refine_topic_bad(state: State) -> dict:
    return {"topic": "优化后的主题", "joke": ""}  # 会把 joke 清空!

6.4 条件边走错路

现象 :图执行到了意料之外的节点。

调试步骤

python

# Step 1: debug 模式查看条件边的路由结果
for chunk in graph.stream(input_data, stream_mode="debug"):
    if chunk["type"] == "task":
        print(f"调度到: {chunk['payload']['name']}")
        print(f"触发条件: {chunk['payload']['triggers']}")

# Step 2: 单独测试路由函数
def route_function(state: State) -> str:
    print(f"路由输入 state: {state}")
    result = "node_a" if state["score"] > 0.5 else "node_b"
    print(f"路由结果: {result}")
    return result

# 在 LangGraph Studio 中可以直接修改 State 重跑
# 验证不同的 State 值会走到哪个分支

七、踩坑记录

坑 1:messages 模式下 Token 不完整

现象 :用 messages 模式拿到最后一个 chunk 时,content 是空的或者不完整。

原因 messages 模式在 LLM 输出结束时会发送一个特殊的结束 chunk,它的 content 为空但 response_metadata 中有完整的 usage 信息。

修复

python

for chunk, metadata in graph.astream(input_data, stream_mode="messages"):
    # 过滤空内容的 chunk
    if not chunk.content:
        continue  # 跳过结束 chunk
    print(chunk.content, end="", flush=True)

坑 2:custom 模式的数据格式问题

现象 get_stream_writer() 推送的数据前端收到后解析失败。

原因 writer() 接受的数据会被序列化为 JSON,但不支持所有 Python 对象(如 datetime、自定义类等)。

修复

python

from datetime import datetime
from langgraph.config import get_stream_writer

def my_node(state):
    writer = get_stream_writer()

    # ❌ 不支持 datetime 直接序列化
    # writer({"time": datetime.now()})

    # ✅ 转为字符串
    writer({"time": datetime.now().isoformat()})

    # ✅ 使用基本类型
    writer({
        "progress": 50,
        "status": "processing",
        "items_done": ["a", "b", "c"],
    })

    return {"result": "done"}

坑 3:流式中断后状态不一致

现象 :流式输出中途被客户端断开,但 Checkpointer 已经保存了部分状态,下次重连时行为异常。

原因 :LangGraph 的 Checkpoint 是在每个 superstep 完成后保存的。如果客户端在节点执行中断开连接, 当前节点的执行会完成 ,Checkpoint 也会保存,但客户端没收到流式输出。

修复

python

# 方案 1:重连后先获取最新状态
config = {"configurable": {"thread_id": thread_id}}
state = graph.get_state(config)

if state.next:  # 还有未执行的节点
    # 从断点恢复
    for chunk in graph.stream(None, config=config, stream_mode="updates"):
        # 重新推送之前没收到的事件
        yield format_sse(chunk)
else:
    # 已经执行完毕,返回最终结果
    yield format_sse({"type": "complete", "data": state.values})

# 方案 2:前端实现幂等消费(根据 step 号去重)

坑 4:astream_events 的事件类型太多难以筛选

现象 astream_events() 返回了几十种事件类型,不知道该关注哪些。

修复 :按需过滤,只关注你关心的事件:

python

# 只关心 LLM 的 token 流和工具调用
IMPORTANT_EVENTS = {
    "on_chat_model_stream",   # LLM Token
    "on_tool_start",          # 工具调用开始
    "on_tool_end",            # 工具调用结束
}

async for event in graph.astream_events(input_data, version="v2"):
    if event["event"] not in IMPORTANT_EVENTS:
        continue

    if event["event"] == "on_chat_model_stream":
        token = event["data"]["chunk"].content
        if token:
            print(token, end="")
    elif event["event"] == "on_tool_start":
        print(f"\n[工具调用] {event['name']}")
    elif event["event"] == "on_tool_end":
        print(f"\n[工具返回] {event['data'].get('output', '')}")

坑 5:LangSmith 采样率过高导致成本飙升

现象 :上线 LangSmith 后,每次请求都在记录 Trace,Token 消耗很大,LangSmith 的费用也跟着涨。

原因 :默认 LANGSMITH_TRACING=true 会记录所有请求。在高并发场景下,Trace 数据量巨大。

修复

python

# 方案 1:采样追踪(只追踪部分请求)
import random
import os

# 10% 的请求开启追踪
if random.random() < 0.1:
    os.environ["LANGSMITH_TRACING"] = "true"
else:
    os.environ["LANGSMITH_TRACING"] = "false"

# 方案 2:按项目/环境分离
# 开发环境全量追踪
os.environ["LANGSMITH_PROJECT"] = "dev-full-trace"

# 生产环境采样追踪
os.environ["LANGSMITH_PROJECT"] = "prod-sampled"

# 方案 3:代码级别控制
from langchain_core.tracers.context import collect_runs

# 只对特定请求开启追踪
result = graph.invoke(input_data, config={"callbacks": []})  # 不追踪
result = graph.invoke(input_data, config={"callbacks": [langsmith_handler]})  # 追踪

总结

plaintext

┌─────────────────────────────────────────────────────────────┐
│               LangGraph 流式 & 调试 速查表                    │
├──────────────┬──────────────────────────────────────────────┤
│              │                                              │
│  流式模式     │  values → 全量 State 快照                    │
│              │  updates → 增量更新                           │
│              │  messages → 逐 Token(ChatBot 必备)          │
│              │  custom → 自定义推送(进度/日志)              │
│              │  debug → 运行时事件                           │
│              │  [组合] → ["messages", "custom"]              │
│              │                                              │
├──────────────┼──────────────────────────────────────────────┤
│              │                                              │
│  调试工具     │  print → 最原始,零配置                       │
│              │  LangGraph Studio → 可视化,回溯,修改重跑     │
│              │  LangSmith Trace → 全链路追踪,Token 统计     │
│              │  自定义回调 → 灵活定制                        │
│              │                                              │
├──────────────┼──────────────────────────────────────────────┤
│              │                                              │
│  异步流式     │  astream() → stream() 的异步版               │
│              │  astream_events() → 更细粒度的事件流          │
│              │  FastAPI + SSE → 生产标准方案                  │
│              │                                              │
├──────────────┼──────────────────────────────────────────────┤
│              │                                              │
│  踩坑注意     │  messages 过滤空 chunk                       │
│              │  custom 数据要可 JSON 序列化                   │
│              │  流式中断注意 Checkpoint 一致性               │
│              │  astream_events 按需过滤事件类型              │
│              │  LangSmith 生产环境建议采样                   │
│              │                                              │
└──────────────┴──────────────────────────────────────────────┘

一句话总结 messages 模式做 ChatBot,updates 模式做进度条,custom 模式做业务推送,debug 模式做开发调试,LangSmith 做生产观测——各司其职,别混着来。

关联阅读

  • 第42篇:LangGraph 快速上手 — 如果你对 LangGraph 还不熟,先看这篇

  • 第41篇:LangGraph 生产部署 — 部署层面也有流式输出的部分内容

  • 第24篇:Agent 可观测性 — 更广义的可观测性话题

0
  1. 支付宝打赏

    qrcode alipay
  2. 微信打赏

    qrcode weixin

评论区