事件流是 LangSmith 部署的类型化投影流模型。LangGraph SDK(Python, JavaScript)针对 LangSmith 部署 API 打开单个订阅,并公开类型化投影——消息、状态、工具调用、子图、输出和自定义转换器扩展——可以从一次运行中并发消费。
事件流位于 流式 API之上,后者公开原始流模式。
快速入门
Python
from langgraph_sdk import get_client
client = get_client(url=DEPLOYMENT_URL, api_key=API_KEY)
async with client.threads.stream(assistant_id="agent") as thread:
await thread.run.start(
input={"messages": [{"role": "user", "content": "What is 42 * 17?"}]},
)
async for message in thread.messages:
async for token in message.text:
print(token, end="", flush=True)
final_state = await thread.output
JavaScript
const client = new Client({
apiUrl: process.env.DEPLOYMENT_URL,
apiKey: process.env.LANGSMITH_API_KEY,
});
const thread = client.threads.stream({ assistantId: "agent" });
await thread.run.start({
input: { messages: [{ role: "user", content: "What is 42 * 17?" }] },
});
for await (const message of thread.messages) {
for await (const token of message.text) {
process.stdout.write(token);
}
}
const finalState = await thread.output;
await thread.close();
JavaScript 流没有 async with 等价物,因此在完成后调用 await thread.close() 释放底层订阅。
cURL
事件流使用两个端点。首先打开 SSE 订阅,然后在同一线程上发送 run.start 命令——SDK 会为你完成这两步,但在网络层面它们是独立的请求。
创建线程:
curl --request POST \
--url /threads \
--header 'Content-Type: application/json' \
--header 'x-api-key: ' \
--data '{}'
打开事件订阅。请求体是一个 EventStreamRequest:要订阅的通道,加上可选的 namespaces, depth和一个 since 序列游标:
curl --request POST \
--url /threads//stream/events \
--header 'Content-Type: application/json' \
--header 'x-api-key: ' \
--data '{"channels": ["values", "updates", "messages", "tools", "lifecycle", "input", "checkpoints", "tasks", "custom"]}'
在第二个请求上发送 run.start 命令以启动运行。命令体是一个 JSON-RPC 风格的信封,包含 id, method和 params:
curl --request POST \
--url /threads//commands \
--header 'Content-Type: application/json' \
--header 'x-api-key: ' \
--data '{
"id": 1,
"method": "run.start",
"params": {
"assistant_id": "agent",
"input": {"messages": [{"role": "user", "content": "What is 42 * 17?"}]}
}
}'
SSE 响应的每一行都是一个 ProtocolEvent 信封;解析事件并按 method 分发以重建 SDK 公开的类型化投影。
有关 LangGraph 应用代码中的进程内等价物,请参阅 LangGraph 事件流.
默认情况下,SDK 通过 Server-Sent Events 进行流式传输。如需使用全双工 WebSocket 连接,请传递 transport="websocket" to client.threads.stream(...).
事件流提供的内容
返回的流通过一个底层事件流公开类型化投影: client.threads.stream(...) 公开类型化投影——
| 投影 | 用途 |
|---|---|
thread.events | 遍历每个原始协议事件(Python)。在 JavaScript 中,打开 thread.subscribe(...). |
thread.messages | 流式传输聊天模型消息、令牌增量、推理和工具调用参数块。 |
thread.values | 遍历状态快照并等待最终值。 |
thread.output | 等待最终输出。 |
thread.tool_calls (thread.toolCalls 在 JavaScript) | 观察具有组装输入、流式输出和结果的工具调用。 |
thread.subgraphs | 发现并观察嵌套图执行。 |
thread.subagents | 面向子代理的视图 thread.subgraphs。在处理 Deep Agents 子代理调用时使用此名称。 |
thread.interrupts | 检查人在环中断的有效载荷。 |
thread.interrupted | 检查运行是否因人工输入而暂停。 |
thread.extensions | 使用发布在以下频道上的自定义流转换器投影。 custom:<name> 频道。 |
多个消费者可以并发读取这些投影。读取 thread.messages 不会消耗 thread.values, thread.toolCalls, thread.subgraphs, or thread.output.
流式消息
使用 thread.messages 处理聊天模型输出:
Python
async with client.threads.stream(assistant_id="agent") as thread:
await thread.run.start(input=input)
async for message in thread.messages:
text = await message.text
usage = (await message.output).usage_metadata
print(text)
print(usage)
JavaScript
const thread = client.threads.stream({ assistantId: "agent" });
await thread.run.start({ input });
for await (const message of thread.messages) {
const text = await message.text;
const usage = (await message.output).usage_metadata;
console.log(text);
console.log(usage);
}
cURL
打开作用域为 messages channel:
curl --request POST \
--url /threads//stream/events \
--header 'Content-Type: application/json' \
--header 'x-api-key: ' \
--data '{"channels": ["messages"]}'
发送 run.start on /commands 如 快速入门中所示。在 params.data.event (message-start, content-block-start, content-block-delta, content-block-finish, message-finish上调度每个事件) 以重新组装每条消息及其内容块。
message.text 既是异步可迭代对象也是可等待对象。迭代它可获得逐 token 输出,等待它可获得完整文本。 message.reasoning 暴露推理增量 message.tool_calls (message.toolCalls 在 JavaScript 中暴露工具调用参数块。等待 message.output 获取最终消息,包括其 usage_metadata。要按精确到达顺序消费文本、推理和工具调用块,请迭代原始事件流而不是分别迭代每个投影。
流式状态
使用 thread.values 在每个步骤后流式传输完整状态快照:
Python
async with client.threads.stream(assistant_id="agent") as thread:
await thread.run.start(input=input)
async for snapshot in thread.values:
print(snapshot)
final_state = await thread.output
JavaScript
const thread = client.threads.stream({ assistantId: "agent" });
await thread.run.start({ input });
for await (const snapshot of thread.values) {
console.log(snapshot);
}
const finalState = await thread.output;
cURL
打开作用域为 values channel:
curl --request POST \
--url /threads//stream/events \
--header 'Content-Type: application/json' \
--header 'x-api-key: ' \
--data '{"channels": ["values"]}'
发送 run.start on /commands 如 快速入门中所示。每个事件的 params.data 是完整状态快照。
thread.values 也是可等待的。等待 thread.values 会解析为最终状态,等同于 await thread.output.
流式工具调用
thread.tool_calls (thread.toolCalls 在 JavaScript 中) 暴露组装后的工具调用。每个句柄包含工具名称 (call.name) 和组装后的输入 (call.input(纯值,未等待)。等待 call.output 以获取调用完成后的工具结果:
Python
async for call in thread.tool_calls:
print(call.name, call.input)
print(await call.output)
if call.error is not None:
print(call.error)
JavaScript
for await (const call of thread.toolCalls) {
console.log(call.name, call.input);
console.log(await call.output);
}
cURL
打开作用域为 tools channel:
curl --request POST \
--url /threads//stream/events \
--header 'Content-Type: application/json' \
--header 'x-api-key: ' \
--data '{"channels": ["tools"]}'
发送 run.start on /commands 如 快速入门中所示。在 params.data.event (tool-started, tool-output-delta, tool-finished, tool-error上调度);通过工具调用 ID 与 messages channel.
在 Python 中, call.deltas 是工具输出流式传输时的异步迭代器,而 call.error 保存工具抛出时的异常。工具事件通过工具调用 ID 与 thread.messages.
流式子图
使用 thread.subgraphs 观察嵌套图工作而无需解析命名空间字符串:
Python
async for subgraph in thread.subgraphs:
print(subgraph.graph_name, subgraph.path)
async for message in subgraph.messages:
print(await message.text)
JavaScript
for await (const subgraph of thread.subgraphs) {
console.log(subgraph.name, subgraph.namespace);
for await (const message of subgraph.messages) {
console.log(await message.text);
}
}
cURL
子图活动通过 params.namespace 路径在每个事件上传达。打开限定于 lifecycle 频道(以及您想在子图内观察的任何频道):
curl --request POST \
--url /threads//stream/events \
--header 'Content-Type: application/json' \
--header 'x-api-key: ' \
--data '{"channels": ["lifecycle", "messages", "tools"]}'
发送 run.start on /commands 如 快速入门中所示。观看 lifecycle 频道以获取 started 事件并 graph_name 来发现新的子图,然后过滤后续事件到该命名空间前缀以观察每个子图的工作。
每个子图句柄公开图名称(subgraph.graph_name 在 Python 中, subgraph.name 在 JavaScript 中)及其命名空间路径(subgraph.path 在 Python 中, subgraph.namespace 在 JavaScript 中),加上每个子图的 messages, tool_calls和嵌套 subgraphs projections.
对于 深度智能体 部署,首选 thread.subagents 用于子智能体调用——它公开子智能体名称以及每个子智能体的消息和工具调用投影。
流式输出
等待 thread.output 运行完成后获取最终状态:
Python
await thread.run.start(input=input)
final_state = await thread.output
JavaScript
await thread.run.start({ input });
const finalState = await thread.output;
cURL
打开限定于 values 和 lifecycle:
curl --request POST \
--url /threads//stream/events \
--header 'Content-Type: application/json' \
--header 'x-api-key: ' \
--data '{"channels": ["values", "lifecycle"]}'
发送 run.start on /commands 如 快速入门中所示。持续读取直到观察到带有 lifecycle 的根命名空间 params.data.event == "completed"事件;最后之前的 values 事件携带最终状态。
thread.output 与 thread.values共享其订阅,因此在读取另一个时等待一个不需要额外的往返。
流式传输多个投影
当应用程序代码需要同时使用多个投影时,运行并发消费者:
Python
async def consume_messages():
async for message in thread.messages:
print(await message.text)
async def consume_tool_calls():
async for call in thread.tool_calls:
print(call.name, await call.output)
async def consume_subgraphs():
async for subgraph in thread.subgraphs:
print(subgraph.graph_name, subgraph.path)
await asyncio.gather(consume_messages(), consume_tool_calls(), consume_subgraphs())
JavaScript
await Promise.all([
(async () => {
for await (const message of thread.messages) {
console.log(await message.text);
}
})(),
(async () => {
for await (const call of thread.toolCalls) {
console.log(call.name, await call.output);
}
})(),
(async () => {
for await (const subgraph of thread.subgraphs) {
console.log(subgraph.name, subgraph.namespace);
}
})(),
]);
cURL
打开一个覆盖您想要消费的所有频道的订阅:
curl --request POST \
--url /threads//stream/events \
--header 'Content-Type: application/json' \
--header 'x-api-key: ' \
--data '{"channels": ["messages", "tools", "lifecycle"]}'
发送 run.start on /commands 如 快速入门中所示。单个 SSE 订阅传递正文中列出的每个频道。根据 method 进行分发以馈送独立消费者;SDK 的并发投影在客户端执行相同的解复用。
每个投影针对同一线程打开过滤订阅,因此并发读取不会增加实际消费的频道之外的服务器负载。
中断后恢复
当图暂停等待人工输入时,检查 thread.interrupted 和 thread.interrupts,然后通过响应中断来恢复:
Python
async with client.threads.stream(assistant_id="agent") as thread:
await thread.run.start(input=input)
async for message in thread.messages:
print(await message.text)
if thread.interrupted:
for interrupt in thread.interrupts:
await thread.run.respond(
{"decisions": [{"type": "approve"}]},
interrupt_id=interrupt["interrupt_id"],
)
final_state = await thread.output
JavaScript
const thread = client.threads.stream({ assistantId: "agent" });
await thread.run.start({ input });
for await (const message of thread.messages) {
console.log(await message.text);
}
if (thread.interrupted) {
for (const interrupt of thread.interrupts) {
await thread.input.respond({
namespace: interrupt.namespace,
interrupt_id: interrupt.interruptId,
response: { decisions: [{ type: "approve" }] },
});
}
}
const finalState = await thread.output;
cURL
将中断响应作为 input.respond command:
curl --request POST \
--url /threads//commands \
--header 'Content-Type: application/json' \
--header 'x-api-key: ' \
--data '{
"id": 2,
"method": "input.respond",
"params": {
"namespace": ,
"interrupt_id": "",
"response": {"decisions": [{"type": "approve"}]}
}
}'
namespace 发送——[] 用于根图);省略它则默认为根。保持原始 SSE 连接打开——部署在命令落地后继续在同一线程上发送事件。
加入一个运行中的任务
要附加到一个已在某线程上运行的任务——在页面重新加载后、在单独的工作线程中或从另一个客户端——使用现有的 thread_id 并跳过 thread.run.start()来打开线程流。部署在连接打开时会重放缓冲的事件,因此消费者可以从头重建任务状态,不会遗漏任何输出。
Python
from langgraph_sdk import get_client
client = get_client(url=DEPLOYMENT_URL, api_key=API_KEY)
async with client.threads.stream(
thread_id=thread_id,
assistant_id="agent",
) as thread:
async for message in thread.messages:
print(await message.text)
final_state = await thread.output
JavaScript
const client = new Client({ apiUrl: DEPLOYMENT_URL, apiKey: API_KEY });
const thread = client.threads.stream(threadId, { assistantId: "agent" });
for await (const message of thread.messages) {
console.log(await message.text);
}
const finalState = await thread.output;
cURL
打开事件流而不发送 run.start 命令来作为被动观察者附加:
curl --request POST \
--url /threads//stream/events \
--header 'Content-Type: application/json' \
--header 'x-api-key: ' \
--data '{"channels": ["values", "updates", "messages", "tools", "lifecycle", "input", "checkpoints", "tasks", "custom"]}'
服务器在订阅打开时从任务开始处重放缓冲的事件。
流式传输所有协议事件
当应用程序代码需要每个事件时,读取原始协议事件流。在 Python 中遍历 thread.events;在 JavaScript 中打开一个 subscribe (JavaScript 流对象本身不可迭代):
Python
async with client.threads.stream(assistant_id="agent") as thread:
await thread.run.start(input=input)
async for event in thread.events:
print(event["method"], event["params"]["namespace"], event["params"]["data"])
要缩小到特定频道,在线程上打开一个 subscribe :
async for event in thread.subscribe(["messages", "tools"]):
...
JavaScript
const thread = client.threads.stream({ assistantId: "agent" });
const events = await thread.subscribe({
channels: ["messages", "tools", "values", "lifecycle"],
});
await thread.run.start({ input });
for await (const event of events) {
console.log(event.method, event.params.namespace, event.params.data);
}
thread.subscribe(...) 返回一个 Promise, so await 在迭代之前使用它。传递一个频道名称数组来缩小订阅范围:
const sub = await thread.subscribe(["messages", "tools"]);
for await (const event of sub) {
// ...
}
cURL
打开一个覆盖所有频道的订阅,然后发送 run.start on /commands ,如 快速入门:
curl --request POST \
--url /threads//stream/events \
--header 'Content-Type: application/json' \
--header 'x-api-key: ' \
--data '{"channels": ["values", "updates", "messages", "tools", "lifecycle", "input", "checkpoints", "tasks", "custom"]}'
响应是一个 SSE 帧流。每个帧有一个 event: 行(频道,映射 method), an id: 行(每个会话的 seq),和一个 data: 行(一个 JSON ProtocolEvent)。持久的 event_id 在 JSON body 内,不在 id: line.
event: lifecycle
id: 1
data: {"type":"event","seq":1,"method":"lifecycle","params":{"namespace":[],"timestamp":1736...,"data":{"event":"started"}},"event_id":"01HZ..."}
event: messages
id: 2
data: {"type":"event","seq":2,"method":"messages","params":{"namespace":[],"timestamp":1736...,"data":{"event":"message-start","message":{...}}},"event_id":"01HZ..."}
每个事件都是一个 ProtocolEvent 包装频道特定载荷的信封:
Python
from typing import Any, NotRequired, TypedDict
class ProtocolEventParams(TypedDict):
namespace: list[str] # path of "<name>:<runtime_id>" segments; [] is the root
timestamp: int # wall-clock milliseconds; can drift, do not rely on for ordering
data: Any # channel-specific payload
class ProtocolEvent(TypedDict):
type: str # always "event"
seq: int # increasing within a session; carried on the SSE `id:` line; use for ordering
method: str # channel name: "messages", "values", "tools", "lifecycle", "custom", ...
params: ProtocolEventParams
event_id: NotRequired[str] # durable cross-session dedup ID, carried in the JSON body
JavaScript
interface ProtocolEvent {
readonly type: "event";
readonly seq: number; // increasing within a session; carried on the SSE `id:` line; use for ordering
readonly method: string; // channel name: "messages", "values", "tools", "lifecycle", "custom", ...
readonly params: {
readonly namespace: string[]; // path of "<name>:<runtime_id>" segments; [] is the root
readonly timestamp: number; // wall-clock milliseconds; can drift, do not rely on for ordering
readonly data: unknown; // channel-specific payload
};
readonly event_id?: string; // durable cross-session dedup ID, carried in the JSON body
}
cURL
原始 SSE 帧:
event: messages
id: 42
data: {"type":"event","seq":42,"method":"messages","params":{"namespace":["researcher:6f4d"],"timestamp":1736283600123,"data":{"event":"content-block-delta","index":0,"delta":{"type":"text","text":"Hello"}}},"event_id":"01HZQ8XK5N6F9M2A3B4C5D6E7F"}
该 id: 行携带每个会话的 seq;持久的 event_id 位于 JSON body 中,是客户端用于去重的关键。线协议不使用 Last-Event-ID 头进行恢复——相反,客户端在请求 body 中传递一个 since 序列游标。Agent Server 在有限缓冲区中按任务缓冲事件,并在新订阅时重放它们,SDK 通过 event_id client-side.
该 namespace 是从根图到发出事件的范围的路径。根是空数组。每个子执行添加一个 "name:runtime_id" 段,因此子图内的嵌套工具调用看起来像 ["researcher:6f4d", "tools:91ac"]。当只需要特定的子树时,直接按命名空间过滤原始事件; thread.subgraphs 已为此对嵌套图执行进行了过滤。
频道和事件生命周期
原始事件在频道上流动。频道名称显示为事件的 method;每个频道发出一个特定的事件形状。
| 频道 | 用途 |
|---|---|
values | 完整图状态快照。 |
updates | 每个节点的状态增量。 |
messages | 以内容块为中心的聊天模型输出。 |
tools | 工具调用开始、流式输出、完成和错误事件。 |
lifecycle | 运行、子图和子代理状态变化。 |
checkpoints | 用于分支和时间旅行的轻量级检查点信封。 |
input | 人工介入输入请求和响应。 |
tasks | Pregel 任务创建和结果事件。 |
custom | 来自图代码的用户定义载荷。 |
custom:<name> | 应用程序定义的流转换器输出。 |
类型化投影(thread.messages, thread.values, thread.toolCalls等)从这些通道构建。通道名称显示为 method 在直接迭代流对象时原始事件上的字段。
消息
messages 通道将输出建模为内容块。 data.event 字段是以下之一 message-start, content-block-start, content-block-delta, content-block-finish, message-finish, or error。内容块具有明确的边界:一个块开始,发送零个或多个增量,并在同一条消息的下一个块开始之前结束。 message-finish 可能包含令牌使用量;不可恢复的模型调用失败作为消息错误事件到达。
工具
tools 通道公开工具执行。 data.event 字段是以下之一 tool-started, tool-output-delta, tool-finished, tool-error。工具事件通过工具调用 ID 关联,因此工具执行可以与 messages channel.
生命周期
lifecycle 通道跟踪根运行、子图和子代理状态。 data.event 字段是以下之一 started, running, completed, failed, interrupted。根运行解析为终端 completed, failed, or interrupted;子作用域也可能报告 running。生命周期数据可能包含可选的 graph_name, error和 cause ,描述子作用域为何启动(父工具调用、扇出发送、边转换)。
从上次事件恢复
事件流是可恢复的。Agent Server 按运行在有界缓冲区中缓冲事件,为每个事件分配一个 seq (按会话排序)和一个持久的 event_id (在重放和副本间稳定),并在重新连接时从游标处重放。SDK 自动处理临时丢失:每个开放订阅跟踪其观察到的最高 seq,在重新连接时 SDK 从该游标重放并通过 event_id.
对重发事件进行去重。要在进程边界之间恢复——页面重新加载、工作进程转移或单独客户端——使用相同的 thread_id重新打开线程。服务器在新订阅打开时重放缓冲的事件,SDK 将其分解到相同的类型化投影中。由于每个运行的缓冲区有界,非常长的运行的最早事件可能已被驱逐。
Python
async with client.threads.stream(
thread_id=thread_id,
assistant_id="agent",
) as thread:
async for event in thread.events:
print(event["method"], event.get("event_id"))
JavaScript
const thread = client.threads.stream(threadId, { assistantId: "agent" });
const events = await thread.subscribe({
channels: ["values", "updates", "messages", "tools", "lifecycle", "input", "checkpoints", "tasks", "custom"],
});
for await (const event of events) {
console.log(event.method, event.event_id);
}
cURL
重新打开订阅以接收重放的事件:
curl --request POST \
--url /threads//stream/events \
--header 'Content-Type: application/json' \
--header 'x-api-key: ' \
--data '{"channels": ["values", "updates", "messages", "tools", "lifecycle", "input", "checkpoints", "tasks", "custom"]}'
在传输层面,请求正文中的 since 序列游标仅返回该点之后的事件。SDK 在重新连接时自动管理此游标,因此 client.threads.stream(...) 不将 since 作为参数暴露。
相关内容
- - 流式 API ——
stream_mode——基于流的 API。同时也支持langgraph-api>=0.10.0. - - LangGraph 事件流 ——相同的概念应用于进程内 LangGraph 应用程序。
- - LangChain 代理事件流 ——针对消息、工具调用和中间件更新的代理聚焦投影。
- - Deep Agents 事件流 ——子代理流、嵌套消息和子代理工具调用。
- - LangSmith 部署 API ——的线路级参考
POST /threads/{thread_id}/stream/events及相关端点。
线路级事件和命令格式定义于 代理协议 repository.