预览版: 尝试对消息、状态、子图、输出和自定义扩展的事件流类型化投影。从事件流概述开始,或在流式输出 cookbook 中探索可运行的示例。
快速开始
基本用法
LangGraph 图暴露了stream(同步)和 astream(异步)方法来以迭代器的形式产出流式输出。传入一个或多个流模式来控制你接收的数据。
for chunk in graph.stream(
{"topic": "ice cream"},
stream_mode=["updates", "custom"],
version="v2",
):
if chunk["type"] == "updates":
for node_name, state in chunk["data"].items():
print(f"节点 {node_name} 已更新:{state}")
elif chunk["type"] == "custom":
print(f"状态:{chunk['data']['status']}")
输出
状态:thinking of a joke...
节点 generate_joke 已更新:{'joke': 'Why did the ice cream go to school? To get a sundae education!'}
完整示例
完整示例
from typing import TypedDict
from langgraph.graph import StateGraph, START, END
from langgraph.config import get_stream_writer
class State(TypedDict):
topic: str
joke: str
def generate_joke(state: State):
writer = get_stream_writer()
writer({"status": "thinking of a joke..."})
return {"joke": f"Why did the {state['topic']} go to school? To get a sundae education!"}
graph = (
StateGraph(State)
.add_node(generate_joke)
.add_edge(START, "generate_joke")
.add_edge("generate_joke", END)
.compile()
)
for chunk in graph.stream(
{"topic": "ice cream"},
stream_mode=["updates", "custom"],
version="v2",
):
if chunk["type"] == "updates":
for node_name, state in chunk["data"].items():
print(f"节点 {node_name} 已更新:{state}")
elif chunk["type"] == "custom":
print(f"状态:{chunk['data']['status']}")
输出
状态:thinking of a joke...
节点 generate_joke 已更新:{'joke': 'Why did the ice cream go to school? To get a sundae education!'}
流输出格式 (v2)
需要 LangGraph >= 1.1。本页所有示例使用
version="v2"。version="v2" 传递给 stream() 或 astream() 以获得统一的输出格式。每个 chunk 都是一个具有一致结构的 StreamPart 字典——无论流模式、模式数量或子图设置如何:
{
"type": "values" | "updates" | "messages" | "custom" | "checkpoints" | "tasks" | "debug",
"ns": (), # 命名空间元组,子图事件时会填充
"data": ..., # 实际负载(类型因流模式而异)
}
TypedDict,包含 ValuesStreamPart、UpdatesStreamPart、MessagesStreamPart、CustomStreamPart、CheckpointStreamPart、TasksStreamPart、DebugStreamPart。你可以从 langgraph.types 中导入这些类型。联合类型 StreamPart 是基于 part["type"] 的不相交联合,可在编辑器和类型检查器中实现完整的类型缩窄。
使用 v1(默认值),输出格式会根据你的流选项变化(单模式返回原始数据,多模式返回 (mode, data) 元组,子图返回 (namespace, data) 元组)。使用 v2,格式始终相同:
for chunk in graph.stream(inputs, stream_mode="updates", version="v2"):
print(chunk["type"]) # "updates"
print(chunk["ns"]) # ()
print(chunk["data"]) # {"node_name": {"key": "value"}}
for chunk in graph.stream(inputs, stream_mode="updates"):
print(chunk) # {"node_name": {"key": "value"}}
chunk["type"] 过滤 chunk 并获得正确的负载类型。每个分支会将 part["data"] 缩窄为该模式的特定类型:
for part in graph.stream(
{"topic": "ice cream"},
stream_mode=["values", "updates", "messages", "custom"],
version="v2",
):
if part["type"] == "values":
# ValuesStreamPart — 每步后的完整状态快照
print(f"状态:topic={part['data']['topic']}")
elif part["type"] == "updates":
# UpdatesStreamPart — 每个节点变更的键
for node_name, state in part["data"].items():
print(f"节点 `{node_name}` 已更新:{state}")
elif part["type"] == "messages":
# MessagesStreamPart — 来自 LLM 调用的 (message_chunk, metadata)
msg, metadata = part["data"]
print(msg.content, end="", flush=True)
elif part["type"] == "custom":
# CustomStreamPart — 来自 get_stream_writer() 的任意数据
print(f"进度:{part['data']['progress']}%")
流模式
将以下一个或多个流模式作为列表传递给stream 或 astream 方法:
| 模式 | 类型 | 描述 |
|---|---|---|
| values | ValuesStreamPart | 每步后的完整状态。 |
| updates | UpdatesStreamPart | 每步后的状态更新。同一步中的多个更新会分别流式输出。 |
| messages | MessagesStreamPart | 来自 LLM 调用的 (LLM Token, 元数据) 二元组。 |
| custom | CustomStreamPart | 通过 get_stream_writer 从节点发出的自定义数据。 |
| checkpoints | CheckpointStreamPart | 检查点事件(与 get_state() 输出格式相同)。需要检查点器。 |
| tasks | TasksStreamPart | 任务开始/完成事件,包含结果和错误。需要检查点器。 |
| debug | DebugStreamPart | 所有可用信息——合并了 checkpoints 和 tasks 并带有额外元数据。 |
图状态
使用流模式updates 和 values 来在图执行时流式输出图的状态。
updates流式输出图每步之后的状态更新。values流式输出图每步之后的完整状态值。
from typing import TypedDict
from langgraph.graph import StateGraph, START, END
class State(TypedDict):
topic: str
joke: str
def refine_topic(state: State):
return {"topic": state["topic"] + " and cats"}
def generate_joke(state: State):
return {"joke": f"This is a joke about {state['topic']}"}
graph = (
StateGraph(State)
.add_node(refine_topic)
.add_node(generate_joke)
.add_edge(START, "refine_topic")
.add_edge("refine_topic", "generate_joke")
.add_edge("generate_joke", END)
.compile()
)
- updates
- values
使用此模式仅流式输出节点在每步之后返回的状态更新。流式输出包含节点名称和更新内容。
for chunk in graph.stream(
{"topic": "ice cream"},
stream_mode="updates",
version="v2",
):
if chunk["type"] == "updates":
for node_name, state in chunk["data"].items():
print(f"节点 `{node_name}` 已更新:{state}")
输出
节点 `refine_topic` 已更新:{'topic': 'ice cream and cats'}
节点 `generate_joke` 已更新:{'joke': 'This is a joke about ice cream and cats'}
使用此模式流式输出图每步之后的完整状态。
for chunk in graph.stream(
{"topic": "ice cream"},
stream_mode="values",
version="v2",
):
if chunk["type"] == "values":
print(f"topic: {chunk['data']['topic']}, joke: {chunk['data']['joke']}")
输出
topic: ice cream, joke:
topic: ice cream and cats, joke:
topic: ice cream and cats, joke: This is a joke about ice cream and cats
LLM Token
使用messages 流模式从图的任何部分(包括节点、工具、子图或任务)逐 Token 流式输出大语言模型(LLM)的输出。
messages 模式的流式输出是一个元组 (message_chunk, metadata),其中:
message_chunk:来自 LLM 的 Token 或消息片段。metadata:包含图节点和 LLM 调用详情的字典。
如果你的 LLM 没有提供 LangChain 集成,你可以使用 custom 模式来流式输出其输出。详见与任意 LLM 配合使用。
Python < 3.11 的异步代码需要手动传递 config
在 Python < 3.11 使用异步代码时,你必须显式将
RunnableConfig 传递给 ainvoke() 以启用正确的流式输出。详见 Python < 3.11 的异步,或升级到 Python 3.11+。from dataclasses import dataclass
from langchain.chat_models import init_chat_model
from langgraph.graph import StateGraph, START
@dataclass
class MyState:
topic: str
joke: str = ""
model = init_chat_model(model="gpt-5.4-mini")
def call_model(state: MyState):
"""调用 LLM 生成关于某主题的笑话"""
# 注意即使使用 .invoke 而非 .stream 运行 LLM,消息事件也会被发出
model_response = model.invoke(
[
{"role": "user", "content": f"Generate a joke about {state.topic}"}
]
)
return {"joke": model_response.content}
graph = (
StateGraph(MyState)
.add_node(call_model)
.add_edge(START, "call_model")
.compile()
)
# "messages" 流模式流式输出 LLM Token 及元数据
# 使用 version="v2" 获得统一的 StreamPart 格式
for chunk in graph.stream(
{"topic": "ice cream"},
stream_mode="messages",
version="v2",
):
if chunk["type"] == "messages":
message_chunk, metadata = chunk["data"]
if message_chunk.content:
print(message_chunk.content, end="|", flush=True)
按 LLM 调用过滤
你可以将tags 与 LLM 调用关联,以按 LLM 调用过滤流式 Token。
from langchain.chat_models import init_chat_model
# model_1 标记为 "joke"
model_1 = init_chat_model(model="gpt-5.4-mini", tags=['joke'])
# model_2 标记为 "poem"
model_2 = init_chat_model(model="gpt-5.4-mini", tags=['poem'])
graph = ... # 定义使用这些 LLM 的图
# stream_mode 设置为 "messages" 以流式输出 LLM Token
# metadata 包含有关 LLM 调用的信息,包括 tags
async for chunk in graph.astream(
{"topic": "cats"},
stream_mode="messages",
version="v2",
):
if chunk["type"] == "messages":
msg, metadata = chunk["data"]
# 通过 metadata 中的 tags 字段过滤流式 Token,
# 仅包含带有 "joke" 标签的 LLM 调用的 Token
if metadata["tags"] == ["joke"]:
print(msg.content, end="|", flush=True)
扩展示例:按标签过滤
扩展示例:按标签过滤
from typing import TypedDict
from langchain.chat_models import init_chat_model
from langgraph.graph import START, StateGraph
# joke_model 标记为 "joke"
joke_model = init_chat_model(model="gpt-5.4-mini", tags=["joke"])
# poem_model 标记为 "poem"
poem_model = init_chat_model(model="gpt-5.4-mini", tags=["poem"])
class State(TypedDict):
topic: str
joke: str
poem: str
async def call_model(state, config):
topic = state["topic"]
print("正在写笑话...")
# 注意:显式传递 config 对于 python < 3.11 是必需的
# 因为之前不支持上下文变量:https://docs.python.org/3/library/asyncio-task.html#creating-tasks
# 显式传递 config 以确保上下文变量正确传播
# 这在 Python < 3.11 使用异步代码时是必需的。请参阅异步部分了解更多详情
joke_response = await joke_model.ainvoke(
[{"role": "user", "content": f"Write a joke about {topic}"}],
config,
)
print("\n\n正在写诗...")
poem_response = await poem_model.ainvoke(
[{"role": "user", "content": f"Write a short poem about {topic}"}],
config,
)
return {"joke": joke_response.content, "poem": poem_response.content}
graph = (
StateGraph(State)
.add_node(call_model)
.add_edge(START, "call_model")
.compile()
)
# stream_mode 设置为 "messages" 以流式输出 LLM Token
# metadata 包含有关 LLM 调用的信息,包括 tags
async for chunk in graph.astream(
{"topic": "cats"},
stream_mode="messages",
version="v2",
):
if chunk["type"] == "messages":
msg, metadata = chunk["data"]
if metadata["tags"] == ["joke"]:
print(msg.content, end="|", flush=True)
从流中排除消息
使用nostream 标签完全从流中排除 LLM 输出。带有 nostream 标签的调用仍然会运行并产生输出;它们的 Token 只是不会在 messages 模式中发出。
这在以下情况下很有用:
- 你需要 LLM 输出用于内部处理(例如结构化输出),但不想将其流式传输给客户端
- 你通过不同的渠道(例如自定义 UI 消息)流式传输相同内容,并希望避免在
messages流中出现重复输出
from typing import Any, TypedDict
from langchain_anthropic import ChatAnthropic
from langgraph.graph import START, StateGraph
stream_model = ChatAnthropic(model_name="claude-haiku-4-5-20251001")
internal_model = ChatAnthropic(model_name="claude-haiku-4-5-20251001").with_config(
{"tags": ["nostream"]}
)
class State(TypedDict):
topic: str
answer: str
notes: str
def answer(state: State) -> dict[str, Any]:
r = stream_model.invoke(
[{"role": "user", "content": f"Reply briefly about {state['topic']}"}]
)
return {"answer": r.content}
def internal_notes(state: State) -> dict[str, Any]:
# Tokens from this model are omitted from stream_mode="messages" because of nostream
r = internal_model.invoke(
[{"role": "user", "content": f"Private notes on {state['topic']}"}]
)
return {"notes": r.content}
graph = (
StateGraph(State)
.add_node("write_answer", answer)
.add_node("internal_notes", internal_notes)
.add_edge(START, "write_answer")
.add_edge("write_answer", "internal_notes")
.compile()
)
initial_state: State = {"topic": "AI", "answer": "", "notes": ""}
stream = graph.stream(initial_state, stream_mode="messages")
按节点过滤
要仅从特定节点流式输出 Token,使用stream_mode="messages" 并通过流式元数据中的 langgraph_node 字段过滤输出:
# "messages" 流模式流式输出 LLM Token 及元数据
# 使用 version="v2" 获得统一的 StreamPart 格式
for chunk in graph.stream(
inputs,
stream_mode="messages",
version="v2",
):
if chunk["type"] == "messages":
msg, metadata = chunk["data"]
# 通过 metadata 中的 langgraph_node 字段过滤流式 Token,
# 仅包含来自指定节点的 Token
if msg.content and metadata["langgraph_node"] == "some_node_name":
...
扩展示例:从特定节点流式输出 LLM Token
扩展示例:从特定节点流式输出 LLM Token
from typing import TypedDict
from langgraph.graph import START, StateGraph
from langchain_openai import ChatOpenAI
model = ChatOpenAI(model="gpt-5.4-mini")
class State(TypedDict):
topic: str
joke: str
poem: str
def write_joke(state: State):
topic = state["topic"]
joke_response = model.invoke(
[{"role": "user", "content": f"Write a joke about {topic}"}]
)
return {"joke": joke_response.content}
def write_poem(state: State):
topic = state["topic"]
poem_response = model.invoke(
[{"role": "user", "content": f"Write a short poem about {topic}"}]
)
return {"poem": poem_response.content}
graph = (
StateGraph(State)
.add_node(write_joke)
.add_node(write_poem)
# 同时写笑话和诗
.add_edge(START, "write_joke")
.add_edge(START, "write_poem")
.compile()
)
# "messages" 流模式流式输出 LLM Token 及元数据
# 使用 version="v2" 获得统一的 StreamPart 格式
for chunk in graph.stream(
{"topic": "cats"},
stream_mode="messages",
version="v2",
):
if chunk["type"] == "messages":
msg, metadata = chunk["data"]
# 通过 metadata 中的 langgraph_node 字段过滤流式 Token,
# 仅包含来自 write_poem 节点的 Token
if msg.content and metadata["langgraph_node"] == "write_poem":
print(msg.content, end="|", flush=True)
自定义数据
要从 LangGraph 节点或工具内部发送自定义用户定义数据,请按照以下步骤操作:- 使用
get_stream_writer访问流写入器并发出自定义数据。 - 调用
.stream()或.astream()时设置stream_mode="custom"以在流中获取自定义数据。你可以组合多个模式(例如["updates", "custom"]),但至少一个必须是"custom"。
Python < 3.11 的异步中无法使用
get_stream_writer
在 Python < 3.11 运行的异步代码中,get_stream_writer 将不起作用。
请改为向你的节点或工具添加 writer 参数并手动传递。
参见 Python < 3.11 的异步了解用法示例。- 节点
- 工具
from typing import TypedDict
from langgraph.config import get_stream_writer
from langgraph.graph import StateGraph, START
class State(TypedDict):
query: str
answer: str
def node(state: State):
# 获取流写入器以发送自定义数据
writer = get_stream_writer()
# 发出自定义键值对(例如进度更新)
writer({"custom_key": "Generating custom data inside node"})
return {"answer": "some data"}
graph = (
StateGraph(State)
.add_node(node)
.add_edge(START, "node")
.compile()
)
inputs = {"query": "example"}
# 设置 stream_mode="custom" 以在流中接收自定义数据
for chunk in graph.stream(inputs, stream_mode="custom", version="v2"):
if chunk["type"] == "custom":
print(f"自定义事件:{chunk['data']['custom_key']}")
from langchain.tools import tool
from langgraph.config import get_stream_writer
@tool
def query_database(query: str) -> str:
"""查询数据库。"""
# 访问流写入器以发送自定义数据
writer = get_stream_writer()
# 发出自定义键值对(例如进度更新)
writer({"data": "Retrieved 0/100 records", "type": "progress"})
# 执行查询
# 发出另一个自定义键值对
writer({"data": "Retrieved 100/100 records", "type": "progress"})
return "some-answer"
graph = ... # 定义使用此工具的图
# 设置 stream_mode="custom" 以在流中接收自定义数据
for chunk in graph.stream(inputs, stream_mode="custom", version="v2"):
if chunk["type"] == "custom":
print(f"{chunk['data']['type']}: {chunk['data']['data']}")
子图输出
要在流式输出中包含来自子图的输出,你可以在父图的.stream() 方法中设置 subgraphs=True。这将流式输出父图和任何子图的输出。
输出将以元组 (namespace, data) 的形式流式传输,其中 namespace 是一个包含调用子图的节点路径的元组,例如 ("parent_node:<task_id>", "child_node:<task_id>")。
- v2 (LangGraph >= 1.1)
- v1(默认)
使用
version="v2" 时,子图事件使用相同的 StreamPart 格式。ns 字段标识来源:for chunk in graph.stream(
{"foo": "foo"},
subgraphs=True,
stream_mode="updates",
version="v2",
):
print(chunk["type"]) # "updates"
print(chunk["ns"]) # () 表示根图,("node_name:<task_id>",) 表示子图
print(chunk["data"]) # {"node_name": {"key": "value"}}
for chunk in graph.stream(
{"foo": "foo"},
# 设置 subgraphs=True 以流式输出子图的输出
subgraphs=True,
stream_mode="updates",
):
print(chunk)
扩展示例:从子图流式输出
扩展示例:从子图流式输出
from langgraph.graph import START, StateGraph
from typing import TypedDict
# 定义子图
class SubgraphState(TypedDict):
foo: str # 注意此键与父图状态共享
bar: str
def subgraph_node_1(state: SubgraphState):
return {"bar": "bar"}
def subgraph_node_2(state: SubgraphState):
return {"foo": state["foo"] + state["bar"]}
subgraph_builder = StateGraph(SubgraphState)
subgraph_builder.add_node(subgraph_node_1)
subgraph_builder.add_node(subgraph_node_2)
subgraph_builder.add_edge(START, "subgraph_node_1")
subgraph_builder.add_edge("subgraph_node_1", "subgraph_node_2")
subgraph = subgraph_builder.compile()
# 定义父图
class ParentState(TypedDict):
foo: str
def node_1(state: ParentState):
return {"foo": "hi! " + state["foo"]}
builder = StateGraph(ParentState)
builder.add_node("node_1", node_1)
builder.add_node("node_2", subgraph)
builder.add_edge(START, "node_1")
builder.add_edge("node_1", "node_2")
graph = builder.compile()
for chunk in graph.stream(
{"foo": "foo"},
stream_mode="updates",
# 设置 subgraphs=True 以流式输出子图的输出
subgraphs=True,
version="v2",
):
if chunk["type"] == "updates":
if chunk["ns"]:
print(f"子图 {chunk['ns']}: {chunk['data']}")
else:
print(f"根图: {chunk['data']}")
根图: {'node_1': {'foo': 'hi! foo'}}
子图 ('node_2:dfddc4ba-c3c5-6887-5012-a243b5b377c2',): {'subgraph_node_1': {'bar': 'bar'}}
子图 ('node_2:dfddc4ba-c3c5-6887-5012-a243b5b377c2',): {'subgraph_node_2': {'foo': 'hi! foobar'}}
根图: {'node_2': {'foo': 'hi! foobar'}}
检查点
使用checkpoints 流模式在图执行时接收检查点事件。每个检查点事件的格式与 get_state() 的输出相同。需要检查点器。
from langgraph.checkpoint.memory import MemorySaver
graph = (
StateGraph(State)
.add_node(refine_topic)
.add_node(generate_joke)
.add_edge(START, "refine_topic")
.add_edge("refine_topic", "generate_joke")
.add_edge("generate_joke", END)
.compile(checkpointer=MemorySaver())
)
config = {"configurable": {"thread_id": "1"}}
for chunk in graph.stream(
{"topic": "ice cream"},
config=config,
stream_mode="checkpoints",
version="v2",
):
if chunk["type"] == "checkpoints":
print(chunk["data"])
任务
使用tasks 流模式在图执行时接收任务开始和完成事件。任务事件包含有关哪个节点正在运行、其结果和任何错误的信息。需要检查点器。
from langgraph.checkpoint.memory import MemorySaver
graph = (
StateGraph(State)
.add_node(refine_topic)
.add_node(generate_joke)
.add_edge(START, "refine_topic")
.add_edge("refine_topic", "generate_joke")
.add_edge("generate_joke", END)
.compile(checkpointer=MemorySaver())
)
config = {"configurable": {"thread_id": "1"}}
for chunk in graph.stream(
{"topic": "ice cream"},
config=config,
stream_mode="tasks",
version="v2",
):
if chunk["type"] == "tasks":
print(chunk["data"])
调试
使用debug 流模式在图执行过程中流式输出尽可能多的信息。流式输出包含节点名称和完整状态。
for chunk in graph.stream(
{"topic": "ice cream"},
stream_mode="debug",
version="v2",
):
if chunk["type"] == "debug":
print(chunk["data"])
debug 模式合并了 checkpoints 和 tasks 事件及额外的元数据。如果你只需要调试信息的子集,请直接使用 checkpoints 或 tasks。同时使用多个模式
你可以将列表作为stream_mode 参数传递,同时流式输出多个模式。
使用 version="v2" 时,每个 chunk 都是一个 StreamPart 字典。使用 chunk["type"] 来区分模式:
for chunk in graph.stream(inputs, stream_mode=["updates", "custom"], version="v2"):
if chunk["type"] == "updates":
for node_name, state in chunk["data"].items():
print(f"节点 `{node_name}` 已更新:{state}")
elif chunk["type"] == "custom":
print(f"自定义事件:{chunk['data']}")
for mode, chunk in graph.stream(inputs, stream_mode=["updates", "custom"]):
print(chunk)
高级
与任意 LLM 配合使用
你可以使用stream_mode="custom" 从任何 LLM API 流式输出数据——即使该 API 没有实现 LangChain 聊天模型接口。
这让你可以集成原始 LLM 客户端或提供自己流式接口的外部服务,使 LangGraph 在自定义设置中非常灵活。
from langgraph.config import get_stream_writer
def call_arbitrary_model(state):
"""调用任意模型并流式输出结果的示例节点"""
# 获取流写入器以发送自定义数据
writer = get_stream_writer()
# 假设你有一个产出 chunk 的流式客户端
# 使用你的自定义流式客户端生成 LLM Token
for chunk in your_custom_streaming_client(state["topic"]):
# 使用 writer 向流发送自定义数据
writer({"custom_llm_chunk": chunk})
return {"result": "completed"}
graph = (
StateGraph(State)
.add_node(call_arbitrary_model)
# 根据需要添加其他节点和边
.compile()
)
# 设置 stream_mode="custom" 以在流中接收自定义数据
for chunk in graph.stream(
{"topic": "cats"},
stream_mode="custom",
version="v2",
):
if chunk["type"] == "custom":
# chunk 数据将包含从 LLM 流式传输的自定义数据
print(chunk["data"])
扩展示例:流式输出任意聊天模型
扩展示例:流式输出任意聊天模型
import operator
import json
from typing import TypedDict
from typing_extensions import Annotated
from langgraph.graph import StateGraph, START
from openai import AsyncOpenAI
openai_client = AsyncOpenAI()
model_name = "gpt-5.4-mini"
async def stream_tokens(model_name: str, messages: list[dict]):
response = await openai_client.chat.completions.create(
messages=messages, model=model_name, stream=True
)
role = None
async for chunk in response:
delta = chunk.choices[0].delta
if delta.role is not None:
role = delta.role
if delta.content:
yield {"role": role, "content": delta.content}
# 这是我们的工具
async def get_items(place: str) -> str:
"""使用此工具列出你被问到的某个地方可能找到的物品。"""
writer = get_stream_writer()
response = ""
async for msg_chunk in stream_tokens(
model_name,
[
{
"role": "user",
"content": (
"Can you tell me what kind of items "
f"i might find in the following place: '{place}'. "
"List at least 3 such items separating them by a comma. "
"And include a brief description of each item."
),
}
],
):
response += msg_chunk["content"]
writer(msg_chunk)
return response
class State(TypedDict):
messages: Annotated[list[dict], operator.add]
# 这是工具调用图节点
async def call_tool(state: State):
ai_message = state["messages"][-1]
tool_call = ai_message["tool_calls"][-1]
function_name = tool_call["function"]["name"]
if function_name != "get_items":
raise ValueError(f"Tool {function_name} not supported")
function_arguments = tool_call["function"]["arguments"]
arguments = json.loads(function_arguments)
function_response = await get_items(**arguments)
tool_message = {
"tool_call_id": tool_call["id"],
"role": "tool",
"name": function_name,
"content": function_response,
}
return {"messages": [tool_message]}
graph = (
StateGraph(State)
.add_node(call_tool)
.add_edge(START, "call_tool")
.compile()
)
AIMessage 来调用图:inputs = {
"messages": [
{
"content": None,
"role": "assistant",
"tool_calls": [
{
"id": "1",
"function": {
"arguments": '{"place":"bedroom"}',
"name": "get_items",
},
"type": "function",
}
],
}
]
}
async for chunk in graph.astream(
inputs,
stream_mode="custom",
version="v2",
):
if chunk["type"] == "custom":
print(chunk["data"]["content"], end="|", flush=True)
禁用特定聊天模型的流式输出
如果你的应用混合使用支持流式输出和不支持流式输出的模型,你可能需要显式禁用不支持流式输出的模型的流式功能。 初始化模型时设置streaming=False。
- init_chat_model
- 聊天模型接口
from langchain.chat_models import init_chat_model
model = init_chat_model(
"claude-sonnet-4-6",
# 设置 streaming=False 以禁用聊天模型的流式输出
streaming=False
)
from langchain_openai import ChatOpenAI
# 设置 streaming=False 以禁用聊天模型的流式输出
model = ChatOpenAI(model="o1-preview", streaming=False)
并非所有聊天模型集成都支持
streaming 参数。如果你的模型不支持,请改用 disable_streaming=True。此参数通过基类在所有聊天模型上可用。迁移到 v2
v2 流式输出格式(本页全程使用)提供了统一的输出格式。以下是主要差异和迁移方法摘要:| 场景 | v1(默认) | v2 (version="v2") |
|---|---|---|
| 单个流模式 | 原始数据(dict) | 带有 type、ns、data 的 StreamPart 字典 |
| 多个流模式 | (mode, data) 元组 | 相同的 StreamPart 字典,按 chunk["type"] 过滤 |
| 子图流式输出 | (namespace, data) 元组 | 相同的 StreamPart 字典,检查 chunk["ns"] |
| 多模式 + 子图 | (namespace, mode, data) 三元组 | 相同的 StreamPart 字典 |
invoke() 返回类型 | 普通 dict(状态) | 带有 .value 和 .interrupts 的 GraphOutput |
| 中断位置(stream) | 状态字典中的 __interrupt__ 键 | values 流部分的 interrupts 字段 |
| 中断位置(invoke) | 结果字典中的 __interrupt__ 键 | GraphOutput 上的 .interrupts 属性 |
| Pydantic/dataclass 输出 | 返回普通 dict | 强制转换为模型/dataclass 实例 |
v2 invoke 格式
当你将version="v2" 传递给 invoke() 或 ainvoke() 时,它返回一个带有 .value 和 .interrupts 属性的 GraphOutput 对象:
from langgraph.types import GraphOutput
result = graph.invoke(inputs, version="v2")
assert isinstance(result, GraphOutput)
result.value # 你的输出——dict、Pydantic 模型或 dataclass
result.interrupts # tuple[Interrupt, ...],如果没有则为空
"values" 以外的任何流模式时,invoke(..., stream_mode="updates", version="v2") 返回 list[StreamPart] 而非 list[tuple]。
GraphOutput 上的字典式访问(result["key"]、"key" in result、result["__interrupt__"])为了向后兼容仍然有效,但已弃用,将在未来版本中移除。请迁移到 result.value 和 result.interrupts。__interrupt__ 下的返回字典中:
config = {"configurable": {"thread_id": "thread-1"}}
result = graph.invoke(inputs, config=config, version="v2")
if result.interrupts:
print(result.interrupts[0].value)
graph.invoke(Command(resume=True), config=config, version="v2")
config = {"configurable": {"thread_id": "thread-1"}}
result = graph.invoke(inputs, config=config)
if "__interrupt__" in result:
print(result["__interrupt__"][0].value)
graph.invoke(Command(resume=True), config=config)
Pydantic 和 dataclass 状态强制转换
当你的图状态是 Pydantic 模型或 dataclass 时,v2 的values 模式会自动将输出强制转换为正确的类型:
from pydantic import BaseModel
from typing import Annotated
import operator
class MyState(BaseModel):
value: str
items: Annotated[list[str], operator.add]
# 使用 version="v2" 时,chunk["data"] 是一个 MyState 实例
for chunk in graph.stream(
{"value": "x", "items": []}, stream_mode="values", version="v2"
):
print(type(chunk["data"])) # <class 'MyState'>
Python < 3.11 的异步
在 Python < 3.11 版本中,asyncio 任务不支持context 参数。
这限制了 LangGraph 自动传播上下文的能力,并在两个关键方面影响了 LangGraph 的流式机制:
- 你必须显式将
RunnableConfig传递给异步 LLM 调用(例如ainvoke()),因为回调不会自动传播。 - 你不能在异步节点或工具中使用
get_stream_writer——你必须直接传递writer参数。
扩展示例:带手动 config 的异步 LLM 调用
扩展示例:带手动 config 的异步 LLM 调用
from typing import TypedDict
from langgraph.graph import START, StateGraph
from langchain.chat_models import init_chat_model
model = init_chat_model(model="gpt-5.4-mini")
class State(TypedDict):
topic: str
joke: str
# 在异步节点函数中接受 config 作为参数
async def call_model(state, config):
topic = state["topic"]
print("正在生成笑话...")
# 将 config 传递给 model.ainvoke() 以确保正确的上下文传播
joke_response = await model.ainvoke(
[{"role": "user", "content": f"Write a joke about {topic}"}],
config,
)
return {"joke": joke_response.content}
graph = (
StateGraph(State)
.add_node(call_model)
.add_edge(START, "call_model")
.compile()
)
# 设置 stream_mode="messages" 以流式输出 LLM Token
async for chunk in graph.astream(
{"topic": "ice cream"},
stream_mode="messages",
version="v2",
):
if chunk["type"] == "messages":
message_chunk, metadata = chunk["data"]
if message_chunk.content:
print(message_chunk.content, end="|", flush=True)
扩展示例:带流写入器的异步自定义流式输出
扩展示例:带流写入器的异步自定义流式输出
from typing import TypedDict
from langgraph.types import StreamWriter
class State(TypedDict):
topic: str
joke: str
# 在异步节点或工具的函数签名中添加 writer 作为参数
# LangGraph 会自动将流写入器传递给函数
async def generate_joke(state: State, writer: StreamWriter):
writer({"custom_key": "Streaming custom data while generating a joke"})
return {"joke": f"This is a joke about {state['topic']}"}
graph = (
StateGraph(State)
.add_node(generate_joke)
.add_edge(START, "generate_joke")
.compile()
)
# 设置 stream_mode="custom" 以在流中接收自定义数据 #
async for chunk in graph.astream(
{"topic": "ice cream"},
stream_mode="custom",
version="v2",
):
if chunk["type"] == "custom":
print(chunk["data"])
将这些文档连接到 Claude、VSCode 等,通过 MCP 获取实时答案。

