Skip to content

本页介绍 Deep Agents 特有的流式关注点——最重要的是通过 stream.subagents 从被委托的子智能体流式输出。有关一般智能体流式输出(stream.messagesstream.values、工具调用、自定义更新),请参阅 LangChain 事件流

流式子智能体

Deep Agents 在 LangGraph 流式输出之上添加了子智能体投影。当你想为每个被委托的 task 调用获取一个流式句柄时,请使用 stream.subagents。该投影是轻量的:它首先发现子智能体任务,而消息、工具调用和值流只在你访问子智能体句柄上的这些投影时才被打开。

每个句柄的 name 是子智能体的配置名称:即协调器在调用 task 工具时传入的 subagent_type。Deep Agents 将该名称绑定到被委托的运行,因此你在子智能体规范中定义的同一标签,就是你在流中用来过滤和路由的标签。

python
stream = agent.stream_events(
    {
        "messages": [{"role": "user", "content": "Write me a haiku about the sea"}],
    },
    version="v3",
)

subagent_names: list[str] = []
for subagent in stream.subagents:
    print(subagent.name, subagent.path, subagent.status)

    for message in subagent.messages:
        print(message.text)

    subagent_names.append(subagent.name)
ts
const stream = await agent.streamEvents(
  { messages: [{ role: "user", content: "Write me a haiku about the sea" }] },
  { version: "v3" },
);

const subagentNames: string[] = [];
for await (const subagent of stream.subagents) {
  console.log(subagent.name);
  console.log(await subagent.taskInput);

  for await (const message of subagent.messages) {
    console.log(await message.text);
  }

  subagentNames.push(subagent.name);
}

子智能体流字段

每个子智能体流都暴露与父级运行相同类型的投影,例如消息、工具调用、嵌套子智能体和最终输出。有关一般父级运行流式模型,请参阅 LangChain 事件流

Python 使用 snake_case 投影名称,例如 tool_calls。每个子智能体流可以暴露 .messages.tool_calls.values.subagents.output

TypeScript 使用 camelCase 投影名称,例如 toolCallstaskInput。每个子智能体流可以暴露 .messages.toolCalls.values.subagents.output

字段描述
name子智能体名称,取自协调器在其 task 调用中选择的 subagent_type
messages子智能体发出的消息。
subagents嵌套的子智能体调用。
output子智能体的最终状态,或委托任务完成的信号。
path子智能体流的命名空间路径。
status生命周期状态,例如 startedcompletedfailedinterrupted
tool_calls作用域限定于子智能体的工具调用。
taskInput传给 task 工具的提示词的 Promise。
toolCalls作用域限定于子智能体的工具调用。

跟踪子智能体生命周期

当你只需要显示哪些子智能体已启动和已完成时,请使用 stream.subagents。除非你在单个子智能体上访问这些投影,否则无需订阅消息或值流。

python
stream = agent.stream_events(input, version="v3")

running = 0
completed = 0
failed = 0

for subagent in stream.subagents:
    running += 1
    print(f"{subagent.name}: started")

    try:
        _ = subagent.output
        running -= 1
        completed += 1
        print(f"{subagent.name}: completed")
    except Exception:
        running -= 1
        failed += 1
        print(f"{subagent.name}: failed")
ts
const stream = await agent.streamEvents(input, { version: "v3" });

let running = 0;
let completed = 0;
let failed = 0;
const watchers: Promise<void>[] = [];

for await (const subagent of stream.subagents) {
  running += 1;
  console.log(`${subagent.name}: started`);

  watchers.push(
    subagent.output.then(
      () => {
        running -= 1;
        completed += 1;
        console.log(`${subagent.name}: completed`);
      },
      () => {
        running -= 1;
        failed += 1;
        console.log(`${subagent.name}: failed`);
      },
    ),
  );
}

await Promise.all(watchers);
console.log({ running, completed, failed });

流式输出消息

Deep Agents 可以从协调器智能体和被委托的子智能体发出消息。对顶层消息使用 stream.messages,对每个被委托的子智能体使用 subagent.messages

python
stream = agent.stream_events(input, version="v3")

coordinator_messages: list[str] = []
for message in stream.messages:
    print("[coordinator]", message.text)
    coordinator_messages.append(message.text)

for subagent in stream.subagents:
    for message in subagent.messages:
        print(f"[{subagent.name}]", message.text)
ts
const stream = await agent.streamEvents(input, { version: "v3" });

const coordinatorMessages: string[] = [];
for await (const message of stream.messages) {
  console.log("[coordinator]", await message.text);
  coordinatorMessages.push(await message.text);
}

for await (const subagent of stream.subagents) {
  for await (const message of subagent.messages) {
    console.log(`[${subagent.name}]`, await message.text);
  }
}

流式输出工具调用

Deep Agents 在智能体树的每一层暴露工具调用。对协调器工具使用顶层 stream.tool_calls,对被委托的工作使用每个 subagent.tool_calls

python
stream = agent.stream_events(input, version="v3")

coordinator_tool_names: list[str] = []
for call in stream.tool_calls:
    print("[coordinator tool]", call.tool_name, call.input)
    print(call.completed, call.error)
    coordinator_tool_names.append(call.tool_name)

for subagent in stream.subagents:
    for call in subagent.tool_calls:
        print(f"[{subagent.name} tool]", call.tool_name, call.input)
        for delta in call.output_deltas:
            print(delta, end="", flush=True)

        if call.completed and call.error is None:
            print(call.output)
        elif call.error is not None:
            print(call.error)
ts
const stream = await agent.streamEvents(input, { version: "v3" });

const coordinatorToolNames: string[] = [];
for await (const call of stream.toolCalls) {
  console.log("[coordinator tool]", call.name, call.input);
  console.log(await call.status);
  coordinatorToolNames.push(call.name);
}

for await (const subagent of stream.subagents) {
  for await (const call of subagent.toolCalls) {
    console.log(`[${subagent.name} tool]`, call.name, call.input);

    const status = await call.status;
    if (status === "finished") {
      console.log(await call.output);
    } else if (status === "error") {
      console.error(await call.error);
    }
  }
}

流式输出嵌套工作

你可以递归进入子智能体流,以观察嵌套的子智能体、消息和工具调用。

python
stream = agent.stream_events(input, version="v3")

subagent_names: list[str] = []
for subagent in stream.subagents:
    print(f"subagent {subagent.name}: {subagent.status}")

    for tool_call in subagent.tool_calls:
        print(f"{tool_call.tool_name}({tool_call.input})")
        for delta in tool_call.output_deltas:
            print(delta, end="", flush=True)

    for nested in subagent.subagents:
        print(f"nested subagent {nested.name}: {nested.status}")

    subagent_names.append(subagent.name)
ts
const stream = await agent.streamEvents(input, { version: "v3" });

const subagentNames: string[] = [];
for await (const subagent of stream.subagents) {
  console.log(`subagent ${subagent.name}: started`);

  for await (const toolCall of subagent.toolCalls) {
    console.log(`${toolCall.name}(${JSON.stringify(toolCall.input)})`);

    const status = await toolCall.status;
    if (status === "finished") {
      console.log(await toolCall.output);
    } else if (status === "error") {
      console.error(await toolCall.error);
    }
  }

  for await (const nested of subagent.subagents) {
    console.log(`nested subagent ${nested.name}: started`);
  }

  subagentNames.push(subagent.name);
}

并发消费

协调器和子智能体的输出经常交错出现。当你需要实时 UI 更新时,请并发消费投影。

对于异步代码中的并发消费,请将 astream_eventsasyncio.gather 一起使用:

py
import asyncio

stream = await agent.astream_events(input, version="v3")

async def consume_coordinator():
    async for message in stream.messages:
        print("[coordinator]", await message.text)

async def consume_subagents():
    async for subagent in stream.subagents:
        async for message in subagent.messages:
            print(f"[{subagent.name}]", await message.text)

await asyncio.gather(consume_coordinator(), consume_subagents())

对于同步代码,请改用 stream.interleave(...)

python
stream = agent.stream_events(input, version="v3")

for name, item in stream.interleave("messages", "subagents"):
    if name == "messages":
        print("[coordinator]", item.text)
    else:
        for message in item.messages:
            print(f"[{item.name}]", message.text)

在 JavaScript 中使用并发消费者:

ts
const stream = await agent.streamEvents(input, { version: "v3" });

await Promise.all([
  (async () => {
    for await (const message of stream.messages) {
      console.log("[coordinator]", await message.text);
    }
  })(),
  (async () => {
    for await (const subagent of stream.subagents) {
      void (async () => {
        for await (const message of subagent.messages) {
          console.log(`[${subagent.name}]`, await message.text);
        }
      })();
    }
  })(),
]);

当你需要跨协调器和所有子智能体的精确到达顺序时,请迭代原始协议事件并使用 namespace 来标识来源:

python
stream = agent.stream_events(input, version="v3")

text_deltas: list[str] = []
for event in stream:
    if event.get("method") != "messages":
        continue

    payload = event["params"]["data"][0]
    if not isinstance(payload, dict):
        continue
    if payload.get("event") != "content-block-delta":
        continue

    block = payload.get("delta") or {}
    if block.get("type") == "text-delta":
        source = "subagent" if event["params"]["namespace"] else "coordinator"
        print(f"[{source}] {block['text']}")
        text_deltas.append(block["text"])
ts
const stream = await agent.streamEvents(input, { version: "v3" });

const textDeltas: string[] = [];
for await (const event of stream) {
  if (event.method !== "messages") continue;

  const data = event.params.data;
  if (data.event !== "content-block-delta") continue;

  const block = data.delta ?? {};
  if (block.type === "text-delta") {
    const isSubagent = event.params.namespace.some((seg) =>
      seg.startsWith("tools:"),
    );
    const source = isSubagent ? "subagent" : "coordinator";
    console.log(`[${source}] ${block.text}`);
    textDeltas.push(block.text);
  }
}

子智能体与子图的对比

stream.subgraphs 显示图的执行结构。stream.subagents 显示产品级的 Deep Agents 任务委托。对于面向用户的 UI,请使用 stream.subagents,因为它隐藏了内部图节点并直接暴露子智能体概念。

相关