以编程方式使用文档

事件流是 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, methodparams:

    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

打开限定于 valueslifecycle:

    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.outputthread.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.interruptedthread.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人工介入输入请求和响应。
tasksPregel 任务创建和结果事件。
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, errorcause ,描述子作用域为何启动(父工具调用、扇出发送、边转换)。

从上次事件恢复

事件流是可恢复的。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.