Skip to content

智能体工作流的数据处理

Supported in ADKPython v2.0.0TypeScript v2.0.0Go v2.0.0

在智能体和基于图的节点之间构建和管理数据,对于使用 ADK 构建可靠的流程至关重要。本指南介绍了基于图的工作流和协作智能体中的数据处理,包括信息如何在图节点之间传输和接收。它涵盖了传递数据、内容和状态的核心参数,并解释了如何使用数据格式 Schema 和特定指令语法为函数节点和智能体节点实现结构化数据传输。

工作流数据流

在基于图的工作流中,节点通过事件向下游步骤传递数据。一个步骤将其输出写入命名的事件字段,下一个步骤将其作为类型化输入接收。

在 Python 中,数据通过 Event 在图节点之间交换。节点数据处理的关键参数包括:

  • output:在节点之间传递信息的参数。
  • message:作为用户回复的数据。
  • state:通过 Event 在整个 ADK 会话中跨节点自动持久化的数据。

在 ADK TypeScript v2.0.0 中,节点通过事件交换数据。节点数据处理的关键字段包括:

  • output:传递给下一个节点的值。直接返回一个值, ADK 会将其包装在事件中,或使用 createEvent({output}) 显式设置该字段。
  • content:面向用户的消息。运行时会渲染此字段, 但图不会将其传递给下一个节点。
  • route:选择要遵循哪条条件边的路由键。

会话状态与事件是分开的。节点通过 ctx.state 读写状态, 累积的增量会附加到该节点的事件上。状态键可以携带前缀来控制其 生命周期和作用域:

前缀 作用域
app: 在应用的所有用户和会话间共享
user: 绑定到用户,在其会话间共享
temp: 当前调用结束后丢弃
(无) 在会话生命周期内持久化

在 ADK Go v2.0.0 中,数据传递机制取决于你使用的智能体风格:

workflow 包(FunctionNode、AgentNode、DynamicNode):节点通过 session.Event 字段进行通信,与 Python 非常相似:

  • Event.Output:节点的返回值,当 FunctionNode 返回非 *genai.Content 值时由框架自动设置。后继节点将其作为类型化 input 参数接收。
  • Event.Routes:由发出节点显式设置的路由键,用于选择要遵循的条件边——相当于 Python 的 Event(route=...)。
  • Event.NodeInfo:调度器元数据(path、MessageAsOutput、OutputFor)。由工作流引擎设置;节点不直接设置此项。

预构建工作流智能体(sequentialagent、parallelagent、loopagent):这些智能体通过会话状态进行通信:

  • llmagent.Config 上的 OutputKey:框架在每轮结束后将智能体的最终文本响应写入 state[OutputKey]。
  • ctx.Session().State().Set / .Get:在自定义代码中对状态进行读写任意值。
  • Instruction 中的 {key}:框架在调用模型之前将 state["key"] 替换到提示词中。

状态键可以携带前缀来控制其生命周期和作用域:

前缀常量 前缀字符串 作用域
session.KeyPrefixApp "app:" 应用中所有用户和会话共享
session.KeyPrefixUser "user:" 绑定到用户,在其会话间共享
session.KeyPrefixTemp "temp:" 当前调用结束后丢弃
(无) — 在会话生命周期内持久化

节点输出

工作流中的每个步骤都会为其后继步骤产生输出。

使用 return 或 yield 语法将数据传递给下一个节点:

from google.adk import Event

def my_function_node(node_input: str):
    output_value = node_input.upper()
    return Event(output=output_value) # "THE RESULT"

当输出不需要额外处理的 Event 数据时,使用 return 语法。当需要发出需要额外处理的数据,或者你正在生成多个数据项时,可以使用多个 yield 命令。每个 yield 调用都会添加到 Event 上的数据对象列表中,该列表会传递给图的下一个节点。不带参数的 return 或 yield 命令会将 None 值传递给下一个节点。

产生节点输出有三种等效方式:直接返回一个值、返回 createEvent({output}), 或从异步生成器中 yield 事件以在结果旁流式传输进度。

import { createEvent, node, NodeContext, Workflow } from '@google/adk';

const returnRawValue = node(
  (_ctx: NodeContext, nodeInput: string) => nodeInput.toUpperCase(),
  { name: 'return_raw_value' },
);

const returnEventOutput = node(
  (_ctx: NodeContext, nodeInput: string) =>
    createEvent({ output: `${nodeInput}!` }),
  { name: 'return_event_output' },
);

const yieldProgressThenOutput = node(
  async function* (_ctx: NodeContext, nodeInput: string) {
    yield createEvent({
      content: { role: 'model', parts: [{ text: 'Working on it...' }] },
    });
    yield createEvent({ output: `<<${nodeInput}>>` });
  },
  { name: 'yield_progress_then_output' },
);

export const rootAgent = new Workflow({
  name: 'node_output_workflow',
  edges: [
    ['START', returnRawValue, returnEventOutput, yieldProgressThenOutput],
  ],
});

注意:每次执行只从一个事件发出 output

一个节点可以 yield 任意数量携带 output 的事件,ADK 在这种情况下不会抛出错误。每个事件会覆盖前一个事件, 后继节点只接收最终值。请改用 content 来发送进度消息。

workflow 包:FunctionNode 只需返回一个类型化的 Go 值。框架会自动将返回值包装在 session.Event 中并设置 Event.Output。后继节点将其作为类型化 input 参数接收——无需手动构建事件:

// newEventOutputPipeline demonstrates the primary data-passing mechanism for
// workflow package nodes: a FunctionNode returns a typed Go value, and the
// framework automatically sets event.Output to that value. The successor node
// receives it as its typed `input` parameter.
//
// This mirrors the Python pattern exactly:
//
//  def my_function_node(node_input: str):
//      return Event(output=node_input.upper())
//
// In Go, the function simply returns the value — no Event construction needed.
func newEventOutputPipeline() (agent.Agent, error) {
    upperFn := func(_ agent.Context, input string) (string, error) {
        return strings.ToUpper(input), nil
    }

    suffixFn := func(_ agent.Context, input string) (string, error) {
        return input + " IS AWESOME!", nil
    }

    nodeA := workflow.NewFunctionNode("upper", upperFn, workflow.NodeConfig{})
    nodeB := workflow.NewFunctionNode("suffix", suffixFn, workflow.NodeConfig{})

    // workflow.Chain wires START → nodeA → nodeB. The output of nodeA is
    // delivered as the typed input of nodeB via event.Output.
    return workflowagent.New(workflowagent.Config{
        Name:        "event_output_pipeline",
        Description: "Demonstrates Event.Output data flow between FunctionNodes.",
        Edges:       workflow.Chain(workflow.Start, nodeA, nodeB),
    })
}

预构建工作流智能体:使用 llmagent.Config 上的 OutputKey 将智能体的文本响应保存到会话状态中,然后在下游智能体的 Instruction 模板中通过 {key} 引用它:

// newOutputKeyPipeline demonstrates the OutputKey mechanism for the prebuilt
// sequentialagent. When OutputKey is set on an llmagent.Config, the framework
// automatically writes the agent's final text response to session state under
// that key. Downstream agents read it by referencing {key} in their Instruction.
//
// This pattern applies to sequentialagent / parallelagent / loopagent.
// For the workflow package (FunctionNode / AgentNode), use Event.Output instead.
func newOutputKeyPipeline(ctx context.Context, geminiModel model.LLM) (agent.Agent, error) {
    step1, err := llmagent.New(llmagent.Config{
        Name:        "step_1",
        Model:       geminiModel,
        Description: "Transforms the user's text.",
        Instruction: "Convert the user's message to uppercase. Output only the transformed text.",
        OutputKey:   "upper_result",
    })
    if err != nil {
        return nil, fmt.Errorf("step1: %w", err)
    }

    step2, err := llmagent.New(llmagent.Config{
        Name:        "step_2",
        Model:       geminiModel,
        Description: "Reports the transformed text.",
        Instruction: "The transformed text is: {upper_result}. Report it to the user.",
    })
    if err != nil {
        return nil, fmt.Errorf("step2: %w", err)
    }

    return sequentialagent.New(sequentialagent.Config{
        AgentConfig: agent.Config{
            Name:      "output_key_pipeline",
            SubAgents: []agent.Agent{step1, step2},
        },
    })
}

节点输出:传递结构化数据

你可以以可序列化的格式传递更长的结构化数据:

def my_function_node_3():
    yield Event(
        output={
            "city_name": "Paris",
            "city_time": "10:10 AM",
        },
    )

注意:Event.output 限制

每次执行只允许节点发出单个 Event.output 数据负载。此限制意味着虽然你可以在一个节点中使用多个 yield,但有两个或更多带有 Event.output 的 yield 命令会导致运行时错误。

output 字段不限于文本。任何可序列化的值都会传递给下一个节点, 下一个节点将其作为类型化对象接收,无需 JSON 解析或状态读取。 在生产节点上附加 outputSchema,或在消费节点上附加 inputSchema, 可以使契约显式化并在运行时进行验证:

import { createEvent, node, NodeContext, Workflow } from '@google/adk';
import { z } from 'zod';

const cityInfoSchema = z.object({
  cityName: z.string(),
  cityTime: z.string(),
});
type CityInfo = z.infer<typeof cityInfoSchema>;

const emitStructuredOutput = node(
  async function* () {
    yield createEvent({
      output: { cityName: 'Paris', cityTime: '10:10 AM' } satisfies CityInfo,
    });
  },
  { name: 'emit_structured_output', outputSchema: cityInfoSchema },
);

const consumeStructuredOutput = node(
  (_ctx: NodeContext, cityInfo: CityInfo) =>
    `It is ${cityInfo.cityTime} in ${cityInfo.cityName} right now.`,
  { name: 'consume_structured_output', inputSchema: cityInfoSchema },
);

export const rootAgent = new Workflow({
  name: 'structured_output_workflow',
  edges: [['START', emitStructuredOutput, consumeStructuredOutput]],
});

workflow 包:FunctionNode 可以返回任何可 JSON 序列化的 Go 结构体。框架将其序列化为 Event.Output,并反序列化为后继节点的类型化 input 参数。没有单个负载限制——每个节点恰好有一个类型化返回值:

// newStructuredOutputPipeline shows how to pass a struct from one FunctionNode
// to another. The framework serialises the return value into event.Output and
// deserialises it back into the successor's typed input parameter.
//
// This is the Go equivalent of:
//
//  class CityTime(BaseModel):
//      time_info: str
//      city: str
//
//  def lookup_time_function(city: str):
//      return Event(output=CityTime(time_info="10:10 AM", city=city))
//
//  def city_report(node_input: CityTime):
//      return Event(output=f"It is {node_input.time_info} in {node_input.city}.")
type CityTime struct {
    TimeInfo string `json:"time_info"`
    City     string `json:"city"`
}

func newStructuredOutputPipeline(ctx context.Context, geminiModel model.LLM) (agent.Agent, error) {
    lookupTimeFn := func(_ agent.Context, city string) (CityTime, error) {
        // Simulate looking up the current time in the city.
        return CityTime{TimeInfo: "10:10 AM", City: city}, nil
    }

    cityReportAgent, err := llmagent.New(llmagent.Config{
        Name:        "city_report_agent",
        Model:       geminiModel,
        Description: "Reports the city and current time from the previous node's output.",
        // When wrapped as an AgentNode, the predecessor's event.Output
        // is delivered as the agent's user content. The {key} template
        // syntax is not required — the struct fields are provided inline.
        Instruction: "Report the city time information you received in a friendly sentence.",
    })
    if err != nil {
        return nil, fmt.Errorf("cityReportAgent: %w", err)
    }

    lookupTimeNode := workflow.NewFunctionNode("lookup_time", lookupTimeFn, workflow.NodeConfig{})
    cityReportNode, err := workflow.NewAgentNode(cityReportAgent, workflow.NodeConfig{})
    if err != nil {
        return nil, fmt.Errorf("NewAgentNode: %w", err)
    }

    return workflowagent.New(workflowagent.Config{
        Name:      "city_time_pipeline",
        Edges:     workflow.Chain(workflow.Start, lookupTimeNode, cityReportNode),
        SubAgents: []agent.Agent{cityReportAgent},
    })
}

预构建工作流智能体:使用多个 OutputKey 值,每个智能体一个,将各个字段存储在会话状态中。下游智能体通过 Instruction 中的 {key} 独立读取每个字段。

路由输出

使用 Event 的 route 参数来驱动条件边分发:

def router(node_input: str):
    return Event(route="BUG")

route 值与 output 独立,因此一个事件可以同时选择分支并向其转发负载。 DEFAULT_ROUTE 设置会捕获没有其他分支匹配的任何值:

import {
  createEvent,
  DEFAULT_ROUTE,
  node,
  NodeContext,
  Workflow,
} from '@google/adk';

const router = node(
  (_ctx: NodeContext, nodeInput: string) =>
    createEvent({
      route: /bug|crash|error/i.test(nodeInput) ? 'BUG' : 'OTHER',
      output: nodeInput,
    }),
  { name: 'router' },
);

const handleBug = node(
  (_ctx: NodeContext, nodeInput: string) => `Filed a bug for: ${nodeInput}`,
  { name: 'handle_bug' },
);

const handleAnythingElse = node(
  (_ctx: NodeContext, nodeInput: string) => `No bug detected in: ${nodeInput}`,
  { name: 'handle_anything_else' },
);

export const rootAgent = new Workflow({
  name: 'routing_output_workflow',
  edges: [
    ['START', router],
    [
      router,
      {
        BUG: handleBug,
        [DEFAULT_ROUTE]: handleAnythingElse,
      },
    ],
  ],
});

workflow 包:发出事件的 FunctionNode 直接构造 session.Event,将 Event.Routes 设置为所需的路由键,并将 Event.Output 设置为将负载转发给后继节点。工作流引擎在分发时读取 Event.Routes 以选择匹配的边:

// classifyAndRoute shows how to set event.Routes alongside event.Output from
// an emitting FunctionNode. The function constructs a session.Event directly,
// sets Routes to select the conditional edge, and sets Output to forward the
// payload to the successor node.
//
// This mirrors the Python pattern:
//
//  def router(node_input: str):
//      return Event(route="BUG")
func classifyAndRoute(ctx agent.Context, msg string, emit func(*session.Event) error) (any, error) {
    category := classifyMessage(msg)

    ev := session.NewEvent(ctx, ctx.InvocationID())
    ev.Routes = []string{category} // drives edge dispatch
    ev.Output = msg                // forwarded as typed input to the successor
    if err := emit(ev); err != nil {
        return nil, err
    }
    return nil, nil // nil suppresses the automatic terminal event
}

func classifyMessage(msg string) string {
    switch {
    case strings.Contains(strings.ToLower(msg), "bug"):
        return "BUG"
    case strings.Contains(strings.ToLower(msg), "help"):
        return "CUSTOMER_SUPPORT"
    default:
        return "LOGISTICS"
    }
}

func newRoutingPipeline() (agent.Agent, error) {
    classifyNode := workflow.NewEmittingFunctionNode("classify", classifyAndRoute, workflow.NodeConfig{})

    bugHandler := workflow.NewFunctionNode("bug_handler",
        func(_ agent.Context, msg string) (string, error) {
            return "Handling bug: " + msg, nil
        }, workflow.NodeConfig{})

    supportHandler := workflow.NewFunctionNode("support_handler",
        func(_ agent.Context, msg string) (string, error) {
            return "Handling support: " + msg, nil
        }, workflow.NodeConfig{})

    logisticsHandler := workflow.NewFunctionNode("logistics_handler",
        func(_ agent.Context, msg string) (string, error) {
            return "Handling logistics: " + msg, nil
        }, workflow.NodeConfig{})

    edges := workflow.Concat(
        workflow.Chain(workflow.Start, classifyNode),
        []workflow.Edge{
            {From: classifyNode, To: bugHandler, Route: workflow.StringRoute("BUG")},
            {From: classifyNode, To: supportHandler, Route: workflow.StringRoute("CUSTOMER_SUPPORT")},
            {From: classifyNode, To: logisticsHandler, Route: workflow.StringRoute("LOGISTICS")},
        },
    )
    return workflowagent.New(workflowagent.Config{
        Name:        "routing_pipeline",
        Description: "Classifies and routes a message using Event.Routes.",
        Edges:       edges,
    })
}

面向用户的消息

使用 Event 的 message 参数向用户发送响应,而不是向下一个节点传递数据:

async def user_message(node_input: str):
  """告知用户研究流程已开始。"""
  yield Event(message="Beginning research process...")

面向用户的消息是事件的 content 字段。运行时会渲染 content, 但图不会将其传递给下一个节点。content 用于面向用户的内容,output 用于传递给下一个节点。一个节点可以通过发送两个事件来同时发出两者, 其中只有一个携带 output:

import { createEvent, node, NodeContext, Workflow } from '@google/adk';

/** Emits a user-facing message: `content`, with no `output`. */
const message = (text: string) =>
  createEvent({ content: { role: 'model', parts: [{ text }] } });

const userMessage = node(
  async function* (_ctx: NodeContext, nodeInput: string) {
    yield message(`Beginning research process for "${nodeInput}"...`);
  },
  { name: 'user_message' },
);

const research = node(
  async function* (_ctx: NodeContext) {
    yield message('Gathering sources...');
    yield createEvent({ output: ['source-a', 'source-b', 'source-c'] });
  },
  { name: 'research' },
);

const report = node(
  (_ctx: NodeContext, sources: string[]) =>
    `Research complete. ${sources.length} sources: ${sources.join(', ')}.`,
  { name: 'report' },
);

export const rootAgent = new Workflow({
  name: 'user_message_workflow',
  edges: [['START', userMessage, research, report]],
});

workflow 包:要在不推进节点类型化输出的情况下发出用户可见的消息,请在通过 EmittingFunctionNode 中的 emit 回调发出的中间事件上设置 Event.Content。最终返回值(或 nil)控制 Event.Output。

预构建工作流智能体:任何 llmagent 步骤都会自动将其模型响应作为面向用户的事件发出。对于非 LLM 步骤,在 agent.Agent 上编写自定义 Run 函数,使其生成 LLMResponse.Content 包含文本的事件。

会话状态和状态作用域

会话状态在会话内的各轮之间持久化数据。它是预构建工作流智能体的主要数据共享机制,无论你使用哪种智能体风格,都可以在工具和回调中使用。

使用 Event 的 state 参数来维护跨节点的值。节点可以修改状态值,修改后的状态值可供下游节点使用:

async def init_state_node(attempts: int = 0):
  yield Event(
      state={
          "attempts": attempts,
      },
  )

async def task_attempt_node(node_input: Content, attempts: int):
  yield Event(
      state={
          "attempts": attempts + 1,
      },
  )

async def read_state_node(ctx: Context):
  print(f"attempts state: {ctx.state}") # attempts state: attempts: 1

root_agent = Workflow(
    name="root_agent",
    edges=[("START", init_state_node, task_attempt_node, read_state_node)],
)

注意:state 属性数据限制

state 参数不应被用于在节点之间持久化大量数据。请使用制品或其他数据持久化机制(如数据库工具)在工作流的生命周期中持久化大型数据资源。

通过 ctx.state 而非返回值来写入状态。写入对同一运行中后续的所有节点可见, 并随写入节点的事件一起提交:

import { node, NodeContext, Workflow } from '@google/adk';

const initStateNode = node(
  (ctx: NodeContext, nodeInput: string) => {
    ctx.state.set('topic', nodeInput.trim());
    ctx.state.set('temp:started_at', new Date().toISOString());
    ctx.state.set('attempts', 0);
  },
  { name: 'init_state_node' },
);

const taskAttemptNode = node(
  (ctx: NodeContext) => {
    const attempts = ctx.state.get<number>('attempts') ?? 0;
    ctx.state.set('attempts', attempts + 1);
  },
  { name: 'task_attempt_node' },
);

const readStateNode = node(
  (ctx: NodeContext) =>
    `attempts state: ${ctx.state.get('attempts')} ` +
    `(topic: ${ctx.state.get('topic')}, ` +
    `started: ${ctx.state.get('temp:started_at')})`,
  { name: 'read_state_node' },
);

export const rootAgent = new Workflow({
  name: 'session_state_workflow',
  edges: [['START', initStateNode, taskAttemptNode, readStateNode]],
});

注意:state 数据限制

会话状态是一个轻量级的键值存储。不要使用它在节点之间传输大型负载; 请改用制品或数据库工具。当只有下一个节点需要某个值时,请将其作为节点 output 沿边传递。当一个值需要在运行结束后继续存在,或需要被工具、 回调或 {key} 指令模板读取时,才使用状态。

状态通过 ctx.Session().State().Set(key, value) 写入,通过 .Get(key) 读取。session 包定义的前缀常量映射到与 Python 的 state 参数相同的生命期作用域。此模式适用于预构建工作流智能体,也适用于任何智能体风格中的工具和回调:

// stateScopes shows how session-state key prefixes control the lifetime and
// visibility of stored values. This pattern applies to the prebuilt workflow
// agents (sequentialagent / parallelagent / loopagent) and to tools and
// callbacks. For the workflow package (FunctionNode / AgentNode), prefer
// returning values directly via Event.Output.
//
// Available prefixes:
//
//  session.KeyPrefixApp  ("app:")  – shared across all users and sessions
//  session.KeyPrefixUser ("user:") – tied to the user, shared across sessions
//  session.KeyPrefixTemp ("temp:") – discarded after the current invocation
//
// Keys with no prefix persist for the lifetime of the session.
func stateScopes(ctx agent.Context) error {
    st := ctx.Session().State()

    // Session-scoped (no prefix) — persists for the life of this session.
    if err := st.Set("attempts", 0); err != nil {
        return fmt.Errorf("state.Set attempts: %w", err)
    }

    // App-scoped — shared across all users and sessions for this app.
    if err := st.Set(session.KeyPrefixApp+"global_counter", 42); err != nil {
        return fmt.Errorf("state.Set app:global_counter: %w", err)
    }

    // User-scoped — shared across all sessions belonging to this user.
    if err := st.Set(session.KeyPrefixUser+"login_count", 1); err != nil {
        return fmt.Errorf("state.Set user:login_count: %w", err)
    }

    // Temp-scoped — discarded after this invocation ends.
    if err := st.Set(session.KeyPrefixTemp+"scratch", "ephemeral"); err != nil {
        return fmt.Errorf("state.Set temp:scratch: %w", err)
    }

    return nil
}

注意:状态数据限制

会话状态是一个轻量级的键值存储。不要使用它来持久化大型负载,如文件内容或二进制数据。请改用 ADK 制品或外部存储工具。

workflow 包:优先使用 Event.Output 而非 state

对于 workflow 包(FunctionNode、AgentNode、DynamicNode),通过返回类型化值在节点之间传递数据——框架会自动设置 Event.Output。只有当你需要与工具、回调或智能体 Instruction 模板共享值时才使用 State().Set。

使用 Schema 约束节点数据

你可以设置输入和输出数据 Schema 来约束任何智能体节点接受和产生的数据格式。

使用扩展自 BaseModel 的类配合 input_schema 和 output_schema 来约束任何智能体的输入和输出:

from google.adk import Agent
from pydantic import BaseModel

class FlightSearchInput(BaseModel):
    origin: str           # 机场代码 "SFO"
    destination: str      # 机场代码 "CDG"
    departure_date: date  # date(2026, 3, 15)
    passengers: int = 1   # 乘客数量

class FlightSearchOutput(BaseModel):
    flights: list[Flight]
    cheapest_price: float

flight_searcher = Agent(
    name="flight_searcher",
    instruction="Search for available flights.",
    input_schema=FlightSearchInput,
    output_schema=FlightSearchOutput,
    tools=[search_flights_api],
    mode="single_turn",
    ...
)

assistant = Agent(
    name="assistant",
    instruction="You help users plan trips.",
    sub_agents=[flight_searcher],
    ...
)

Schema 是 Zod 对象或 genai Schema。Schema 的位置决定其效果:

  • LlmAgent.outputSchema 选项要求模型以该形状回答。
  • LlmAgent.inputSchema 选项仅在智能体作为工具暴露时适用。 在图内部,使用 node(agent, {inputSchema}) 在节点本身上设置 验证节点输入的 Schema。

图中的智能体必须以 single_turn 模式运行(这是默认值),或 task 模式。

import {
  FunctionTool,
  LlmAgent,
  node,
  NodeContext,
  Workflow,
} from '@google/adk';
import { z } from 'zod';

const flightSearchInputSchema = z.object({
  origin: z.string().describe('Origin airport code, e.g. "SFO".'),
  destination: z.string().describe('Destination airport code, e.g. "CDG".'),
  departureDate: z.string().describe('Departure date, e.g. "2026-03-15".'),
  passengers: z.number().describe('Number of passengers.'),
});
type FlightSearchInput = z.infer<typeof flightSearchInputSchema>;

const flightSchema = z.object({
  carrier: z.string(),
  flightNumber: z.string(),
  price: z.number(),
});

const flightSearchOutputSchema = z.object({
  flights: z.array(flightSchema),
  cheapestPrice: z.number(),
});
type FlightSearchOutput = z.infer<typeof flightSearchOutputSchema>;

/** Stands in for a real flight-search API. */
const searchFlightsApi = new FunctionTool({
  name: 'search_flights_api',
  description: 'Searches available flights for a route and date.',
  parameters: flightSearchInputSchema,
  execute: ({ origin, destination }) => [
    {
      carrier: 'AF',
      flightNumber: `AF${origin.length}${destination.length}0`,
      price: 812.4,
    },
    {
      carrier: 'UA',
      flightNumber: `UA${origin.length}${destination.length}1`,
      price: 947.0,
    },
  ],
});

const parseRequest = node(
  (_ctx: NodeContext, nodeInput: string): FlightSearchInput => {
    const codes = nodeInput.toUpperCase().match(/\b[A-Z]{3}\b/g) ?? [];
    const date = nodeInput.match(/\d{4}-\d{2}-\d{2}/)?.[0];
    const passengers = Number(
      nodeInput.match(/(\d+)\s*(people|pax|passengers?)/i)?.[1],
    );
    return {
      origin: codes[0] ?? 'SFO',
      destination: codes[1] ?? 'CDG',
      departureDate: date ?? '2026-03-15',
      passengers: Number.isFinite(passengers) ? passengers : 1,
    };
  },
  { name: 'parse_request', outputSchema: flightSearchInputSchema },
);

const flightSearcher = new LlmAgent({
  name: 'flight_searcher',
  model: 'gemini-flash-latest',
  mode: 'single_turn',
  instruction:
    'Search for available flights with the search_flights_api tool and report ' +
    'every flight it returns plus the cheapest price.',
  inputSchema: flightSearchInputSchema,
  outputSchema: flightSearchOutputSchema,
  tools: [searchFlightsApi],
});

const renderResults = node(
  (_ctx: NodeContext, results: FlightSearchOutput) =>
    `Cheapest: $${results.cheapestPrice}\n` +
    results.flights
      .map((f) => `  ${f.carrier} ${f.flightNumber} — $${f.price}`)
      .join('\n'),
  { name: 'render_results', inputSchema: flightSearchOutputSchema },
);

export const rootAgent = new Workflow({
  name: 'flight_workflow',
  edges: [
    [
      'START',
      parseRequest,
      node(flightSearcher, { inputSchema: flightSearchInputSchema }),
      renderResults,
    ],
  ],
});

workflow 包:使用 workflow.NewAgentNodeTyped[Input, Output] 为智能体节点附加 Schema。泛型类型参数会自动反射为 *jsonschema.Schema——无需手动构建 Schema。节点的 Event.Output 将结构化结果传递给后继节点——不需要 OutputKey 或状态写入:

// FlightSearchInput is the typed input schema for the flight-search agent node.
// workflow.NewAgentNodeTyped[FlightSearchInput, FlightSearchOutput] reflects
// these structs into *jsonschema.Schema automatically — no hand-built schema
// construction needed.
type FlightSearchInput struct {
    Origin        string `json:"origin"         jsonschema:"Departure airport code e.g. SFO"`
    Destination   string `json:"destination"    jsonschema:"Arrival airport code e.g. CDG"`
    DepartureDate string `json:"departure_date" jsonschema:"Travel date in YYYY-MM-DD format"`
}

// FlightSearchOutput is the typed output schema for the flight-search agent node.
type FlightSearchOutput struct {
    CheapestPrice string `json:"cheapest_price" jsonschema:"Cheapest available fare e.g. $450"`
    FlightCount   string `json:"flight_count"   jsonschema:"Number of matching flights found"`
}

// newSchemaAgentPipeline demonstrates workflow.NewAgentNodeTyped, which infers
// *jsonschema.Schema from the generic type parameters. This is the Go equivalent
// of Python's:
//
//  flight_searcher = Agent(
//      input_schema=FlightSearchInput,
//      output_schema=FlightSearchOutput,
//      ...
//  )
//
// The node's event.Output carries the structured result to the successor —
// no OutputKey or state write is needed.
func newSchemaAgentPipeline(ctx context.Context, geminiModel model.LLM) (agent.Agent, error) {
    flightSearchAgent, err := llmagent.New(llmagent.Config{
        Name:        "flight_searcher",
        Model:       geminiModel,
        Description: "Searches for available flights and returns structured results.",
        Instruction: `You are a flight-search assistant. Respond ONLY with a JSON object.`,
    })
    if err != nil {
        return nil, fmt.Errorf("flightSearchAgent: %w", err)
    }

    synthAgent, err := llmagent.New(llmagent.Config{
        Name:        "trip_assistant",
        Model:       geminiModel,
        Description: "Summarises flight search results for the user.",
        Instruction: `You help users plan trips. Summarise the flight result you received.`,
    })
    if err != nil {
        return nil, fmt.Errorf("synthAgent: %w", err)
    }

    // NewAgentNodeTyped[In, Out] reflects FlightSearchInput and FlightSearchOutput
    // into *jsonschema.Schema automatically. The node enforces the input schema
    // and constrains the model reply to the output schema's shape.
    flightNode, err := workflow.NewAgentNodeTyped[FlightSearchInput, FlightSearchOutput](flightSearchAgent, workflow.NodeConfig{})
    if err != nil {
        return nil, fmt.Errorf("flightNode: %w", err)
    }

    synthNode, err := workflow.NewAgentNode(synthAgent, workflow.NodeConfig{})
    if err != nil {
        return nil, fmt.Errorf("synthNode: %w", err)
    }

    return workflowagent.New(workflowagent.Config{
        Name:      "flight_booking_pipeline",
        Edges:     workflow.Chain(workflow.Start, flightNode, synthNode),
        SubAgents: []agent.Agent{flightSearchAgent, synthAgent},
    })
}

预构建工作流智能体:在 llmagent.Config 上设置 InputSchema 和 OutputSchema。OutputSchema 强制模型回复符合 Schema 的 JSON 对象(当设置了 OutputSchema 时智能体无法使用工具)。使用 OutputKey 将 JSON 字符串保存到状态中,供下游智能体通过 Instruction 中的 {key} 引用。

在智能体中访问结构化数据

使用花括号 { } 语法从输入 Schema 中选择属性,或使用 < > 选择属性并通过源节点名称进行限定:

class CityTime(BaseModel):
    time_info: str  # 时间信息
    city: str       # 城市名称

def lookup_time_function(city: str):
    """模拟返回指定城市的当前时间。"""
    return Event(output=CityTime(time_info='10:10 AM', city=city))

city_report_agent = Agent(
    name="city_report_agent",
    model="gemini-flash-latest",
    input_schema=CityTime,

    # 基于类和参数的数据选择
    # instruction="""
    #     Return a sentence in the following format:
    #     It is {CityTime.time_info} in {CityTime.city} right now.
    # """,

    # 基于源节点名称的更严格数据选择
    instruction="""
        Return a sentence in the following format:
        It is <CityTime.time_info from lookup_time_function> in
        <CityTime.city from lookup_time_function> right now.
    """,
)

root_agent = Workflow(
    name="root_agent",
    edges=[
        (START, city_generator_agent, lookup_time_function, city_report_agent)
    ],
)

在智能体指令中有两种数据选择形式:

  • {Class.field} 形式从当前节点的输入中读取字段。
  • <Class.field from source_node> 形式从指定前驱节点的输出中读取字段。 当多个上游节点共享相同字段名时使用此形式。

两种形式都与 {state_key} 不同,后者读取会话状态。 Class. 前缀仅用于文档说明;解析时使用点号后的字段名。

import { LlmAgent, node, NodeContext, Workflow } from '@google/adk';
import { z } from 'zod';

const cityTimeSchema = z.object({
  timeInfo: z.string().describe('Time information.'),
  city: z.string().describe('City name.'),
});
type CityTime = z.infer<typeof cityTimeSchema>;

const cityGeneratorAgent = new LlmAgent({
  name: 'city_generator_agent',
  model: 'gemini-flash-latest',
  instruction: 'Return the name of a random city. Return only the name.',
});

/** Simulates returning the current time in the specified city. */
const lookupTimeFunction = node(
  (_ctx: NodeContext, city: string): CityTime => ({
    timeInfo: '10:10 AM',
    city: city.trim(),
  }),
  { name: 'lookup_time_function', outputSchema: cityTimeSchema },
);

const cityReportAgent = new LlmAgent({
  name: 'city_report_agent',
  model: 'gemini-flash-latest',
  instruction:
    'Return a sentence in the following format: It is ' +
    '<CityTime.timeInfo from lookup_time_function> in ' +
    '<CityTime.city from lookup_time_function> right now.',
});

export const rootAgent = new Workflow({
  name: 'root_agent',
  edges: [
    [
      'START',
      cityGeneratorAgent,
      lookupTimeFunction,
      node(cityReportAgent, { inputSchema: cityTimeSchema }),
    ],
  ],
});

在 ADK Go v2.0.0 中,FunctionNode 返回一个类型化结构体,框架将其序列化为 Event.Output。后继的 AgentNode 将该结构体作为用户内容接收——字段可直接用于智能体的 Instruction,无需任何 {key} 模板语法。这相当于 Python 的 input_schema=CityTime 配合 {CityTime.time_info} 模板占位符:结构化字段作为类型化输入传递,而非从状态中按名称查找。

// newStructuredOutputPipeline shows how to pass a struct from one FunctionNode
// to another. The framework serialises the return value into event.Output and
// deserialises it back into the successor's typed input parameter.
//
// This is the Go equivalent of:
//
//  class CityTime(BaseModel):
//      time_info: str
//      city: str
//
//  def lookup_time_function(city: str):
//      return Event(output=CityTime(time_info="10:10 AM", city=city))
//
//  def city_report(node_input: CityTime):
//      return Event(output=f"It is {node_input.time_info} in {node_input.city}.")
type CityTime struct {
    TimeInfo string `json:"time_info"`
    City     string `json:"city"`
}

func newStructuredOutputPipeline(ctx context.Context, geminiModel model.LLM) (agent.Agent, error) {
    lookupTimeFn := func(_ agent.Context, city string) (CityTime, error) {
        // Simulate looking up the current time in the city.
        return CityTime{TimeInfo: "10:10 AM", City: city}, nil
    }

    cityReportAgent, err := llmagent.New(llmagent.Config{
        Name:        "city_report_agent",
        Model:       geminiModel,
        Description: "Reports the city and current time from the previous node's output.",
        // When wrapped as an AgentNode, the predecessor's event.Output
        // is delivered as the agent's user content. The {key} template
        // syntax is not required — the struct fields are provided inline.
        Instruction: "Report the city time information you received in a friendly sentence.",
    })
    if err != nil {
        return nil, fmt.Errorf("cityReportAgent: %w", err)
    }

    lookupTimeNode := workflow.NewFunctionNode("lookup_time", lookupTimeFn, workflow.NodeConfig{})
    cityReportNode, err := workflow.NewAgentNode(cityReportAgent, workflow.NodeConfig{})
    if err != nil {
        return nil, fmt.Errorf("NewAgentNode: %w", err)
    }

    return workflowagent.New(workflowagent.Config{
        Name:      "city_time_pipeline",
        Edges:     workflow.Chain(workflow.Start, lookupTimeNode, cityReportNode),
        SubAgents: []agent.Agent{cityReportAgent},
    })
}

有关此工作流的完整示例,请参阅基于图的智能体工作流。