动态智能体工作流¶
ADK 框架提供了一种编程方式来定义工作流,作为基于图的工作流的更灵活、更强大的替代方案。使用基于图的方法可以方便地通过工作流节点组合多步骤的静态流程结构。然而,如果你的工作流逻辑路径更复杂,包含迭代循环或复杂的分支逻辑,基于图的方法可能不适合你的需求,或者可能变得过于笨重而难以管理。
ADK 中的动态工作流允许你抛开基于图的路径结构,使用所选编程语言的全部能力来构建工作流。通过动态工作流,你可以使用简单的装饰器(Python)或构造函数(Go)创建工作流,将工作流节点作为函数调用,并构建复杂的路由逻辑。以下是 ADK 动态工作流的一些优势:
- 灵活的控制流: 使用循环、条件判断和递归来动态定义执行顺序,这些在静态图中很难或无法表示。
- 编程体验: 使用熟悉的构造,如
while循环和async/await(Python)或for循环和workflow.RunNode(Go),而不是基于图的路由。 - 自动检查点: 动态工作流会跟踪每个节点的执行。恢复工作流时会自动跳过已成功的子节点,使复杂逻辑默认具有持久性和可恢复性。
- 封装: 将业务逻辑包装到父节点中,在内部组合低级节点,使整体工作流保持清晰和可管理。
开始使用¶
以下动态工作流代码示例展示了如何定义一个包含单个节点和函数的基本工作流:
from google.adk import Context
from google.adk import Workflow
from google.adk.workflow import node
from typing import Any
@node(name="hello_node")
def my_node(node_input: Any):
return "Hello World"
# 定义一个动态工作流节点
@node(rerun_on_resume=True)
async def my_workflow(ctx: Context, node_input: str) -> str:
# run_node 执行一个节点并返回其输出
result = await ctx.run_node(my_node, node_input="hello")
return result
# 运行工作流
root_agent = Workflow(
name="root_agent",
edges=[("START", my_workflow)],
)
此示例使用 @node 注解以简化代码,保持代码尽可能简洁。此注解会生成包装器,使代码可以在 ADK 动态工作流的上下文中运行。
TypeScript 没有 @node 装饰器。请改用 node(fn, options) 工厂
函数。ctx.runNode() 方法等同于 ctx.run_node():
import { node, NodeContext, Workflow } from '@google/adk';
const myNode = node(() => 'Hello World', { name: 'hello_node' });
const myWorkflow = node(
async (ctx: NodeContext, _nodeInput: string) => {
const result = await ctx.runNode(myNode, 'hello');
return result.output;
},
{ name: 'my_workflow', rerunOnResume: true },
);
export const rootAgent = new Workflow({
name: 'root_agent',
edges: [['START', myWorkflow]],
});
当你编写编排器节点时,两个细节会影响你读取结果的方式以及 工作流在暂停后的行为:
ctx.runNode()方法解析为节点结果,而非输出值。 读取.output属性以获取值。- 调用
ctx.runNode()的编排器必须设置rerunOnResume: true。此设置会导致节点主体在恢复时重新运行, 已完成的子节点会从其检查点重放,而不会再次执行。
在 Go 中,workflow.NewFunctionNode 替代了 @node 装饰器,workflow.NewDynamicNode 替代了 @node(rerun_on_resume=True) 异步编排器。workflow.RunNode 等同于 ctx.run_node()。使用 workflowagent.New 和 workflow.Chain 替代 Workflow(edges=[...])。
人工介入暂停后的恢复行为由 NodeConfig.RerunOnResume 控制——详情请参见下方的节点。
// helloNode is a simple FunctionNode that returns "Hello World".
// In Python this would be written as:
//
// @node(name="hello_node")
// def my_node(node_input: Any):
// return "Hello World"
//
// In Go, workflow.NewFunctionNode wraps the same logic with the
// required node interface, inferring input and output types from
// the generic parameters.
var helloNode = workflow.NewFunctionNode("hello_node",
func(_ agent.Context, _ string) (string, error) {
return "Hello World", nil
},
workflow.NodeConfig{},
)
// myWorkflow is a dynamic orchestrator node. It calls workflow.RunNode
// to schedule helloNode as a child and returns its output.
// In Python this would be:
//
// @node(rerun_on_resume=True)
// async def my_workflow(ctx: Context, node_input: str) -> str:
// result = await ctx.run_node(my_node, node_input="hello")
// return result
//
// workflow.NewDynamicNode defaults RerunOnResume to &true, matching the
// Python @node(rerun_on_resume=True) behaviour.
var myWorkflow = workflow.NewDynamicNode[string, string]("my_workflow",
func(ctx agent.Context, _ string, _ func(*session.Event) error) (string, error) {
return workflow.RunNode[string](ctx, helloNode, "hello")
},
workflow.NodeConfig{},
)
func runGetStarted() error {
ctx := context.Background()
// workflowagent.New creates an agent.Agent backed by the workflow engine.
// workflow.Chain(workflow.Start, myWorkflow) produces the edges slice
// equivalent to Python's edges=[("START", my_workflow)].
wa, err := workflowagent.New(workflowagent.Config{
Name: "root_agent",
Description: "A minimal dynamic workflow.",
Edges: workflow.Chain(workflow.Start, myWorkflow),
})
if err != nil {
return fmt.Errorf("workflowagent.New: %w", err)
}
l := full.NewLauncher()
return l.Execute(ctx, &launcher.Config{
AgentLoader: agent.NewSingleLoader(wa),
}, os.Args[1:])
}
构建块:节点和工作流¶
节点和工作流是 ADK 动态工作流的基本构建块。这些类型和函数提供了所需的功能,可以包装你的代码,使其能够集成到 ADK 基于代码的工作流中。
Nodes¶
ADK 中的动态工作流由节点组成。一个简单的工作流节点包装了一个普通函数,并附带在工作流中运行所需的元数据。
在 Python 中,@node 注解会生成节点包装器,将样板代码降到最低:
以下代码片段展示了不使用 @node 注解的等效代码:
# 基础函数
def my_function_node(node_input: Any):
return "Hello World"
# 带选项的 FunctionNode 包装器
success_node = FunctionNode(
my_function_node,
name="hello",
rerun_on_resume=True,
)
手动创建节点包装器代码在以下情况会很有用:当你要包装来自外部库的函数时,需要从同一函数创建具有不同配置的多个节点时,或者当你要在注册表中管理节点引用以进行高级编排时。
有两种方式来构建节点:node(fn, options) 工厂函数,和显式的
new FunctionNode(name, fn, config) 构造函数。当你包装来自其他库的函数、
需要从同一函数创建多个不同配置的节点,或在注册表中管理节点引用以进行
高级编排时,使用构造函数。
import { FunctionNode, node, NodeContext, Workflow } from '@google/adk';
/** The plain function both node forms wrap. */
function myFunctionNode(_ctx: NodeContext, nodeInput: unknown): string {
return `Hello ${nodeInput ?? 'World'}`;
}
const helloNode = node(myFunctionNode, { name: 'hello_node' });
const successNode = new FunctionNode('hello', myFunctionNode, {
rerunOnResume: true,
});
在此代码示例中,最重要的选项是 rerunOnResume,它控制工作流在
人工在回路暂停后恢复时的行为:
true(重新进入): 节点主体从头重新运行。对任何调用ctx.runNode()的编排器使用此设置。主体会重新执行, 已完成的子激活会自动跳过。false(交接,叶子节点的默认值): 恢复负载被路由到节点的 后继节点作为输入,绕过被中断的节点。
在 Go 中,workflow.NewFunctionNode[IN, OUT] 将普通函数包装为工作流节点,并从泛型参数推断输入和输出类型。没有装饰器语法;节点是一个值,你需要将其作为子节点传递给动态编排器中的 workflow.RunNode:
// myFunctionNode demonstrates the explicit NewFunctionNode constructor —
// equivalent to wrapping a function in a FunctionNode manually in Python:
//
// success_node = FunctionNode(my_function_node, name="hello", rerun_on_resume=True)
//
// Creating the node directly (rather than via @node) is useful when you
// need multiple nodes from the same function with different configurations,
// or when wrapping functions from an external library.
var myFunctionNode = workflow.NewFunctionNode("hello",
func(_ agent.Context, _ any) (string, error) {
return "Hello World", nil
},
workflow.NodeConfig{},
)
// myFormattingNode is a second function node that the dynamic orchestrator
// calls in sequence, mirroring:
//
// result_formatted = await ctx.run_node(my_formatting_node, node_input=result)
var myFormattingNode = workflow.NewFunctionNode("format",
func(_ agent.Context, in string) (string, error) {
return fmt.Sprintf("[formatted] %s", in), nil
},
workflow.NodeConfig{},
)
NodeConfig 与 Python 的 @node 参数持有相同的选项。最重要的字段是 RerunOnResume *bool,它控制工作流在人工介入暂停后恢复时的行为:
&true(重新进入模式):恢复时从头重新运行被中断的节点。适用于在循环中调用workflow.RunNode的动态编排器节点——主体会重新执行,已完成的子激活会自动跳过(检查点)。这与 Python 的@node(rerun_on_resume=True)对应。&false(交接模式):恢复时将 payload 直接路由到节点的后继节点作为输入,完全绕过被中断的节点。适用于只发出暂停事件并期望人工响应流向下一步的叶子节点。nil:默认行为取决于节点类型。workflow.NewDynamicNode自动将nil → &true(重新进入模式),因为编排器主体必须在恢复时重新进入以传递缓存的子结果。workflow.NewFunctionNode和其他叶子节点构造函数保持nil不变,引擎将其视为交接(&false)。在任何节点类型上,显式的&false始终会被尊重。
// NewDynamicNode: nil RerunOnResume 自动设置为 &true。
// 显式传递 &rerun 是等效的,且意图更清晰。
rerun := true
orchestratorNode := workflow.NewDynamicNode[string, string]("my_workflow",
myOrchestratorfn,
workflow.NodeConfig{RerunOnResume: &rerun}, // 重新进入:节点主体在恢复时重新运行
)
// NewFunctionNode: nil RerunOnResume 保持 nil → 引擎将其视为交接。
handoffNode := workflow.NewFunctionNode("leaf_node",
myLeafFn,
workflow.NodeConfig{}, // nil RerunOnResume → FunctionNode 的交接模式
)
Workflows¶
在 ADK 动态工作流中,你使用动态节点作为节点的主要编排器。动态节点管理子节点的运行以及这些节点的执行逻辑(顺序和路径)。
@node(rerun_on_resume=True)
async def my_workflow(ctx):
# run_node 执行一个节点并返回其输出
result = await ctx.run_node(my_function_node, node_input="Hello")
result_formatted = await ctx.run_node(my_formatting_node, node_input=result)
return result_formatted
# 运行工作流
root_agent = Workflow(
name="root_agent",
edges=[("START", my_workflow)],
)
编排器是一个异步函数,为每个子步骤 await ctx.runNode()。
使用 rerunOnResume: true 将其包装为节点,并将其作为图的唯一边:
const myFormattingNode = node(
(_ctx: NodeContext, nodeInput: string) => `>> ${nodeInput.trim()} <<`,
{ name: 'my_formatting_node' },
);
const myWorkflow = node(
async (ctx: NodeContext, nodeInput: unknown) => {
const greeted = await ctx.runNode(helloNode, nodeInput);
const again = await ctx.runNode(successNode, greeted.output);
const formatted = await ctx.runNode(myFormattingNode, again.output);
return formatted.output;
},
{ name: 'my_workflow', rerunOnResume: true },
);
export const rootAgent = new Workflow({
name: 'root_agent',
edges: [['START', myWorkflow]],
});
workflow.NewDynamicNode 创建一个编排器,其主体为每个子步骤调用 workflow.RunNode。使用 workflowagent.New 和 workflow.Chain(workflow.Start, myWorkflow) 等同于 Workflow(edges=[("START", my_workflow)]):
// orchestratorWorkflow is a dynamic node that schedules two children in
// sequence via workflow.RunNode, equivalent to:
//
// @node(rerun_on_resume=True)
// async def my_workflow(ctx):
// result = await ctx.run_node(my_function_node, node_input="Hello")
// result_formatted = await ctx.run_node(my_formatting_node, node_input=result)
// return result_formatted
var orchestratorWorkflow = workflow.NewDynamicNode[string, string]("my_workflow",
func(ctx agent.Context, _ string, _ func(*session.Event) error) (string, error) {
result, err := workflow.RunNode[string](ctx, myFunctionNode, "Hello")
if err != nil {
return "", err
}
return workflow.RunNode[string](ctx, myFormattingNode, result)
},
workflow.NodeConfig{},
)
数据处理¶
在使用 ADK 动态工作流时,传递数据比基于图的工作流更简单,因为 workflow.RunNode 直接以类型化的 Go 值返回子节点的输出——消除了手动读写会话状态键来进行数据传输的需要。
from google.adk import Context
from google.adk.workflow import node
@node(rerun_on_resume=True)
async def editorial_workflow(ctx: Context, user_request: str):
# 智能体节点生成输出
raw_draft = await ctx.run_node(draft_agent, user_request)
# 函数节点格式化文本
formatted_text = await ctx.run_node(format_function_node, raw_draft)
return formatted_text
你还可以使用定义的类传递特定的数据模式,并配置输入和输出模式,类似于基于图的工作流节点:
from google.adk import Agent
from google.adk import Context
from google.adk.workflow import node
from pydantic import BaseModel
class CityTime(BaseModel):
time_info: str # 时间信息
city: str # 城市名称
@node
def city_time_function(city: str):
"""模拟返回指定城市的当前时间。"""
return CityTime(time_info="10:10 AM", city=city)
city_report_agent = Agent(
name="city_report_agent",
model="gemini-flash-latest",
input_schema=CityTime,
instruction="""output the data provided by the previous node.""",
)
@node # 工作流节点
async def city_workflow(ctx: Context):
city_time = await ctx.run_node(city_time_function, "Paris")
report_text = await ctx.run_node(city_report_agent, city_time)
return report_text
ctx.runNode() 函数直接返回子节点的结果,因此无需读写会话状态键
即可将值向下游传递一步。此函数接受任何类节点值,包括 LlmAgent,
无需先将其包装在 node() 中:
import { LlmAgent, node, NodeContext, Workflow } from '@google/adk';
const draftAgent = new LlmAgent({
name: 'draft_agent',
model: 'gemini-flash-latest',
instruction: 'Write a short draft for the user request.',
});
const formatFunctionNode = node(
(_ctx: NodeContext, rawDraft: string) =>
rawDraft
.split('\n')
.map((line) => line.trim())
.filter(Boolean)
.map((line) => `| ${line}`)
.join('\n'),
{ name: 'format_function_node' },
);
const editorialWorkflow = node(
async (ctx: NodeContext, userRequest: string) => {
const rawDraft = await ctx.runNode(draftAgent, userRequest);
const formattedText = await ctx.runNode(
formatFunctionNode,
rawDraft.output,
);
return formattedText.output;
},
{ name: 'editorial_workflow', rerunOnResume: true },
);
export const rootAgent = new Workflow({
name: 'root_agent',
edges: [['START', editorialWorkflow]],
});
Schema 的工作方式与图中的相同。将其附加到你运行的节点上, 如顺序路由部分所示。
在 Go 中,workflow.NewAgentNode 包装一个 agent.Agent,使其可以通过动态编排器中的 workflow.RunNode 调用。每个 RunNode 调用的输出以类型化的值返回——不需要读取会话状态:
// newDataHandlingWorkflow demonstrates how to pass data between a dynamic
// orchestrator and an LlmAgent-backed node. workflow.NewAgentNode wraps an
// agent.Agent so it can be invoked via workflow.RunNode.
//
// In Python this mirrors:
//
// city_report_agent = Agent(name="city_report_agent", ...)
// @node
// async def city_workflow(ctx: Context):
// city_time = await ctx.run_node(city_time_function, "Paris")
// report_text = await ctx.run_node(city_report_agent, city_time)
// return report_text
func newDataHandlingWorkflow(ctx context.Context) (agent.Agent, error) {
model, err := gemini.NewModel(ctx, "gemini-flash-latest", &genai.ClientConfig{})
if err != nil {
return nil, fmt.Errorf("gemini.NewModel: %w", err)
}
// cityTimeNode is a FunctionNode that returns a formatted city-time string.
cityTimeNode := workflow.NewFunctionNode("city_time_function",
func(_ agent.Context, city string) (string, error) {
return fmt.Sprintf("10:10 AM in %s", city), nil
},
workflow.NodeConfig{},
)
// cityReportAgent is an LlmAgent that receives the city-time string and
// produces a human-friendly report.
cityReportAgent, err := llmagent.New(llmagent.Config{
Name: "city_report_agent",
Model: model,
Description: "Reports city time information.",
Instruction: "Output the data provided by the previous node in a friendly sentence.",
})
if err != nil {
return nil, fmt.Errorf("llmagent.New (cityReport): %w", err)
}
// workflow.NewAgentNode wraps cityReportAgent so it can be called from
// inside a dynamic node via workflow.RunNode.
cityReportNode, err := workflow.NewAgentNode(cityReportAgent, workflow.NodeConfig{})
if err != nil {
return nil, fmt.Errorf("workflow.NewAgentNode: %w", err)
}
cityWorkflow := workflow.NewDynamicNode[string, string]("city_workflow",
func(ctx agent.Context, _ string, _ func(*session.Event) error) (string, error) {
cityTime, err := workflow.RunNode[string](ctx, cityTimeNode, "Paris")
if err != nil {
return "", err
}
return workflow.RunNode[string](ctx, cityReportNode, cityTime)
},
workflow.NodeConfig{},
)
return workflowagent.New(workflowagent.Config{
Name: "data_handling_workflow",
SubAgents: []agent.Agent{cityReportAgent},
Edges: workflow.Chain(workflow.Start, cityWorkflow),
})
}
有关工作流节点之间数据处理的更多信息,请参见智能体工作流的数据处理。
工作流路由¶
与基于图的工作流相比,ADK 中的动态工作流在路由逻辑方面提供了更大的灵活性,包括迭代循环或更复杂的分支逻辑。本节描述了一些你可以使用的路由技术。
Sequence route¶
与基于图的工作流一样,你可以使用 ADK 动态工作流创建顺序任务处理。
以下代码片段展示了一个动态工作流,包含一个智能体、一个函数节点和第二个智能体:
顺序路由依次等待 ctx.runNode() 调用。每个调用在下一个开始前完成:
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 a specified city. */
const cityTimeFunction = node(
(_ctx: NodeContext, city: string): CityTime => ({
timeInfo: '10:10 AM',
city: city.trim(),
}),
{ name: 'city_time_function', outputSchema: cityTimeSchema },
);
const cityReportAgent = node(
new LlmAgent({
name: 'city_report_agent',
model: 'gemini-flash-latest',
instruction: 'Output the data provided by the previous node as a sentence.',
}),
{ inputSchema: cityTimeSchema },
);
const cityWorkflow = node(
async (ctx: NodeContext) => {
const city = await ctx.runNode(cityGeneratorAgent);
const cityTime = await ctx.runNode(cityTimeFunction, city.output);
const reportText = await ctx.runNode(cityReportAgent, cityTime.output);
return reportText.output;
},
{ name: 'city_workflow', rerunOnResume: true },
);
export const rootAgent = new Workflow({
name: 'root_agent',
edges: [['START', cityWorkflow]],
});
在 NewDynamicNode 主体中顺序调用 workflow.RunNode——每个调用会等待子节点完成后再开始下一个。上面的数据处理示例恰好展示了这种模式:cityWorkflow 按顺序调用 workflow.RunNode 处理 cityTimeNode,然后是 cityReportNode,将每个节点的类型化输出传递给下一个。
Loop route¶
对于你想使用迭代循环来处理任务的工作流,动态工作流在定义所需路由逻辑方面提供了更大的灵活性。
以下代码示例展示了如何使用动态工作流构建用于生成、审查和更新代码的工作流循环:
from google.adk import Context
from google.adk import Event
from google.adk.agents import LlmAgent
from google.adk.workflow import node
coder_agent = LlmAgent(
name="generator_agent",
model="gemini-flash-latest",
instruction="Write python code for user request.",
)
@node(name="lint_reviewer")
async def compile_lint_check(ctx: Context, code: str):
# 模拟 API 调用或 lint 检查
class Response:
findings = ""
return Response()
fixer_agent = LlmAgent(
name="fixer_agent",
model="gemini-flash-latest",
instruction="""Refactor current code {code}.
Based on compile & lint review: {findings}""",
)
@node # 工作流节点
async def code_workflow(ctx: Context, user_request: str):
code = await ctx.run_node(coder_agent, user_request)
check_resp = await ctx.run_node(compile_lint_check, code)
while check_resp.findings:
yield Event(state={"code": code, "findings": check_resp.findings})
code = await ctx.run_node(fixer_agent, {"code": code, "findings": check_resp.findings})
check_resp = await ctx.run_node(compile_lint_check, code)
yield Event(output=code)
动态工作流通过将迭代定义为普通循环而非图中的回边,有助于保持工作流逻辑简洁。 值保存在局部变量中,状态仅在智能体指令模板需要读回时才写入。 与图循环不同,循环受其循环条件约束:
import { LlmAgent, node, NodeContext, Workflow } from '@google/adk';
/** Safety bound on the refine loop. */
const MAX_FIX_ROUNDS = 3;
const coderAgent = new LlmAgent({
name: 'generator_agent',
model: 'gemini-flash-latest',
instruction: 'Write TypeScript code for the user request. Output code only.',
});
/** Simulates a compile / lint pass. Empty findings means "clean". */
const compileLintCheck = node(
(_ctx: NodeContext, code: string) => {
const findings: string[] = [];
if (!/\/\*\*/.test(code)) {
findings.push('every function needs a JSDoc comment');
}
if (!/\)\s*:\s*\w/.test(code)) {
findings.push('add return type annotations');
}
return { findings: findings.join('; ') };
},
{ name: 'lint_reviewer' },
);
const fixerAgent = new LlmAgent({
name: 'fixer_agent',
model: 'gemini-flash-latest',
instruction: `Refactor current code {code}.
Based on compile & lint review: {findings}
Output code only.`,
});
const codeWorkflow = node(
async (ctx: NodeContext, userRequest: string) => {
let code = (await ctx.runNode(coderAgent, userRequest)).output as string;
let checkResp = (await ctx.runNode(compileLintCheck, code)).output as {
findings: string;
};
for (let round = 0; checkResp.findings && round < MAX_FIX_ROUNDS; round++) {
ctx.state.set('code', code);
ctx.state.set('findings', checkResp.findings);
code = (
await ctx.runNode(fixerAgent, { code, findings: checkResp.findings })
).output as string;
checkResp = (await ctx.runNode(compileLintCheck, code)).output as {
findings: string;
};
}
return code;
},
{ name: 'code_workflow', rerunOnResume: true },
);
export const rootAgent = new Workflow({
name: 'root_agent',
edges: [['START', codeWorkflow]],
});
在 Go 中,循环是动态节点主体中的普通 for 循环。当没有发现时,lint 检查节点返回空字符串,信号循环退出:
// newLoopWorkflow demonstrates an iterative loop inside a dynamic node.
// The orchestrator body uses a plain Go for loop to keep calling the
// lintCheckNode until there are no findings — equivalent to Python's:
//
// @node
// async def code_workflow(ctx: Context, user_request: str):
// code = await ctx.run_node(coder_agent, user_request)
// check_resp = await ctx.run_node(compile_lint_check, code)
// while check_resp.findings:
// code = await ctx.run_node(fixer_agent, ...)
// check_resp = await ctx.run_node(compile_lint_check, code)
// return code
func newLoopWorkflow(ctx context.Context) (agent.Agent, error) {
model, err := gemini.NewModel(ctx, "gemini-flash-latest", &genai.ClientConfig{})
if err != nil {
return nil, fmt.Errorf("gemini.NewModel: %w", err)
}
coderAgent, err := llmagent.New(llmagent.Config{
Name: "generator_agent",
Model: model,
Description: "Writes Go code for the user request.",
Instruction: "Write Go code for the user request. Output only the code.",
OutputKey: "generated_code",
})
if err != nil {
return nil, fmt.Errorf("llmagent.New (coder): %w", err)
}
coderNode, err := workflow.NewAgentNode(coderAgent, workflow.NodeConfig{})
if err != nil {
return nil, fmt.Errorf("workflow.NewAgentNode (coder): %w", err)
}
// lintCheckNode simulates a lint/compile check. It returns an empty
// string when there are no findings, signalling the loop to exit.
lintCheckNode := workflow.NewFunctionNode("lint_reviewer",
func(_ agent.Context, code string) (string, error) {
// Simulate a lint check: return findings or empty string when clean.
if len(code) < 50 {
return "Code is too short; add error handling.", nil
}
return "", nil // no findings — loop exits
},
workflow.NodeConfig{},
)
fixerAgent, err := llmagent.New(llmagent.Config{
Name: "fixer_agent",
Model: model,
Description: "Refactors code based on lint findings.",
Instruction: "Refactor the provided code to address the review findings. Output only the improved code.",
})
if err != nil {
return nil, fmt.Errorf("llmagent.New (fixer): %w", err)
}
fixerNode, err := workflow.NewAgentNode(fixerAgent, workflow.NodeConfig{})
if err != nil {
return nil, fmt.Errorf("workflow.NewAgentNode (fixer): %w", err)
}
codeWorkflow := workflow.NewDynamicNode[string, string]("code_workflow",
func(ctx agent.Context, userRequest string, _ func(*session.Event) error) (string, error) {
code, err := workflow.RunNode[string](ctx, coderNode, userRequest)
if err != nil {
return "", err
}
findings, err := workflow.RunNode[string](ctx, lintCheckNode, code)
if err != nil {
return "", err
}
// Loop until the lint check reports no findings.
for findings != "" {
code, err = workflow.RunNode[string](ctx, fixerNode, code)
if err != nil {
return "", err
}
findings, err = workflow.RunNode[string](ctx, lintCheckNode, code)
if err != nil {
return "", err
}
}
return code, nil
},
workflow.NodeConfig{},
)
return workflowagent.New(workflowagent.Config{
Name: "code_pipeline",
SubAgents: []agent.Agent{coderAgent, fixerAgent},
Edges: workflow.Chain(workflow.Start, codeWorkflow),
})
}
Parallel execution routes¶
ADK 中的动态工作流可以支持并行执行。
在 Python 中,你可以使用 asyncio.gather 来构建并行执行:
import asyncio
from typing import Any
from google.adk import Context
from google.adk.workflow import BaseNode, node
@node(rerun_on_resume=True)
async def parallel_supervisor(
ctx: Context, node_input: list[Any], real_node: BaseNode
):
"""并行运行工作节点,处理输入列表中的每个项。"""
tasks = []
for item in node_input:
# ctx.run_node 返回一个 future。追加而不是立即等待。
tasks.append(ctx.run_node(real_node, item))
# 并行收集所有结果
results = await asyncio.gather(*tasks)
return results
提示:恢复并行节点
工作流框架确保如果动态工作流被恢复,只有失败或中断的工作节点会被重新执行,包括并行工作节点。
ctx.runNode() 方法返回一个 Promise,因此在等待任何子节点之前启动所有子节点
会并发运行子节点,Promise.all 收集结果。运行 ID 按调用顺序分配,
因此在同步循环中启动子节点以保持 ID 在恢复时的确定性:
import { node, NodeContext, Workflow } from '@google/adk';
const sleep = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms));
/** The worker run once per list item. */
const realNode = node(
async (_ctx: NodeContext, item: string) => {
await sleep(200);
return { item, length: item.length };
},
{ name: 'analyze_item' },
);
const parallelSupervisor = node(
async (ctx: NodeContext, nodeInput: string) => {
const items = nodeInput
.split(',')
.map((item) => item.trim())
.filter(Boolean);
const tasks = items.map((item) => ctx.runNode(realNode, item));
const results = await Promise.all(tasks);
return results.map((result) => result.output);
},
{ name: 'parallel_supervisor', rerunOnResume: true },
);
const summarize = node(
(_ctx: NodeContext, results: Array<{ item: string; length: number }>) =>
results.map((r) => `${r.item}: ${r.length} chars`).join('\n'),
{ name: 'summarize' },
);
export const rootAgent = new Workflow({
name: 'root_agent',
edges: [['START', parallelSupervisor, summarize]],
});
提示:优先使用内置的并行工作者
要对列表中的每个项运行同一个节点,请使用
node(worker, {parallelWorker: true, maxParallelWorkers: 4})。
此选项执行扇出并限制并发数(默认为 8)。当你需要自定义调度
或部分失败处理时,使用上面展示的手动方式。在恢复时,
两种方式中只有失败或中断的工作者才会重新执行。
在 Go 中,workflow.NewParallelWorker 包装一个子节点,并对列表输入的每个元素并发运行它,将结果收集到单个输出切片中。maxConcurrency 参数限制同时运行的并发激活数量;0 表示无限制:
// newParallelWorkflow demonstrates parallel execution using
// workflow.NewParallelWorker. The worker node runs a wrapped child node
// concurrently for each element in a list input, collecting results.
//
// This is the Go equivalent of using asyncio.gather in Python:
//
// @node(rerun_on_resume=True)
// async def parallel_supervisor(ctx, node_input, real_node):
// tasks = [ctx.run_node(real_node, item) for item in node_input]
// results = await asyncio.gather(*tasks)
// return results
func newParallelWorkflow() (agent.Agent, error) {
// workerNode processes a single item. NewParallelWorker will call it
// once per element of the list input, concurrently.
workerNode := workflow.NewFunctionNode("worker",
func(_ agent.Context, item string) (string, error) {
return fmt.Sprintf("processed: %s", item), nil
},
workflow.NodeConfig{},
)
// NewParallelWorker wraps workerNode so it runs concurrently for each
// element of a []string input. maxConcurrency=0 means unlimited.
parallelWorker, err := workflow.NewParallelWorker(
"parallel_supervisor",
workerNode,
0, // maxConcurrency: 0 = unlimited
workflow.NodeConfig{},
)
if err != nil {
return nil, fmt.Errorf("workflow.NewParallelWorker: %w", err)
}
return workflowagent.New(workflowagent.Config{
Name: "parallel_workflow",
Description: "Runs a worker node in parallel for each item in the input list.",
Edges: workflow.Chain(workflow.Start, parallelWorker),
})
}
提示:恢复并行节点
工作流框架确保如果动态工作流被恢复,只有失败或中断的工作节点会被重新执行,包括由 NewParallelWorker 管理的并行工作节点。
人工输入¶
ADK 中的动态工作流还可以包含人工输入或人工在回路(HITL)步骤。
你可以通过从节点生成 RequestInput 来将人工输入构建到工作流中,这会暂停工作流并等待用户输入。以下代码示例展示了如何构建人工输入节点并将其包含在工作流中:
from typing import Any
from google.adk import Context
from google.adk.events import RequestInput
from google.adk.workflow import node
@node(rerun_on_resume=False)
async def get_user_approval(ctx: Context, node_input: Any):
"""生成 RequestInput 以暂停工作流并等待用户输入。"""
yield RequestInput(message="Please approve this request (Yes/No)")
@node(rerun_on_resume=True)
async def handle_process(ctx: Context, node_input: Any):
"""编排器调用交互式步骤。"""
user_response = await ctx.run_node(get_user_approval)
if user_response.lower() == "yes":
return "Approved"
return "Denied"
重要:使用 ctx.run_node 的父节点
动态工作流中调用 ctx.run_node 的父节点必须设置 rerun_on_resume=True 以正确处理中断。
叶子节点返回 RequestInput 以暂停工作流,并保持默认的 rerunOnResume: false,
使回复成为其输出。调用它的编排器必须设置 rerunOnResume: true:
import { node, NodeContext, RequestInput, Workflow } from '@google/adk';
/**
* Pauses the workflow and waits for user input.
*
* `rerunOnResume: false` (the default, spelled out here because it is the
* point) is what makes this a one-liner: the reply is handed to the node as
* its output instead of the body running a second time to collect it.
*/
const getUserApproval = node(
() => new RequestInput({ message: 'Please approve this request (Yes/No)' }),
{ name: 'get_user_approval', rerunOnResume: false },
);
/** The orchestrator calling the interactive step. */
const handleProcess = node(
async (ctx: NodeContext, nodeInput: unknown) => {
const approval = await ctx.runNode(getUserApproval, nodeInput);
if (approval.interruptIds.length > 0) {
return undefined;
}
const userResponse = String(approval.output ?? '')
.trim()
.toLowerCase();
if (userResponse === 'yes') {
return 'Approved';
}
return 'Denied';
},
{ name: 'handle_process', rerunOnResume: true },
);
export const rootAgent = new Workflow({
name: 'root_agent',
edges: [['START', handleProcess]],
});
重要:在做决定前检查 interruptIds
ctx.runNode() 方法在子节点中断时不会抛出错误。
它正常返回,结果的 interruptIds 属性被填充,而 output
属性仍为 undefined。在使用结果之前检查 interruptIds。
跳过此检查的编排器会将缺失的输出视为答案,
并使用用户从未提供的值继续执行。
在 Go 中,使用 workflow.NewEmittingFunctionNode 和 workflow.ResumeOrRequestInput 来实现重新进入的 HITL 模式。在第一次通过时,ResumeOrRequestInput 发出 session.RequestInput 事件并返回 ErrNodeInterrupted,暂停工作流。人工回复后,节点从头重新运行(RerunOnResume: &true),ResumeOrRequestInput 直接返回人工的回复:
// newHITLWorkflow demonstrates the re-entry HITL pattern using
// workflow.ResumeOrRequestInput. On the first pass the node emits a
// RequestInput event and returns ErrNodeInterrupted (pausing the workflow).
// After the human replies, the same node is re-run from the top
// (RerunOnResume=&true) and ResumeOrRequestInput returns the human's reply.
//
// In Python this is equivalent to:
//
// @node(rerun_on_resume=True)
// async def get_user_approval(ctx, node_input):
// yield RequestInput(message="Please approve this request (Yes/No)")
//
// @node(rerun_on_resume=True)
// async def handle_process(ctx, node_input):
// user_response = await ctx.run_node(get_user_approval)
// if user_response.lower() == "yes":
// return "Approved"
// return "Denied"
func newHITLWorkflow() (agent.Agent, error) {
rerun := true
// approvalNode pauses on the first pass to ask the user for a Yes/No
// approval, then resolves their decision on resume.
// workflow.ResumeOrRequestInput handles both phases.
approvalNode := workflow.NewEmittingFunctionNode[any, any]("get_user_approval",
func(nc agent.Context, _ any, emit func(*session.Event) error) (any, error) {
// ResumeOrRequestInput: on first pass, emits the prompt and
// returns ErrNodeInterrupted. On re-run after the human replies,
// it returns the reply payload directly.
reply, err := workflow.ResumeOrRequestInput(nc, emit, session.RequestInput{
InterruptID: "user_approval",
Message: "Please approve this request (Yes/No)",
})
if err != nil {
return nil, err
}
response, _ := reply.(string)
if response == "" {
response = "No"
}
if response == "yes" || response == "Yes" {
return "Approved", nil
}
return "Denied", nil
},
workflow.NodeConfig{RerunOnResume: &rerun},
)
return workflowagent.New(workflowagent.Config{
Name: "hitl_workflow",
Description: "Pauses for user approval before completing a task.",
Edges: workflow.Chain(workflow.Start, approvalNode),
})
}
高级功能¶
动态工作流提供了一些旨在处理更复杂开发场景的高级功能。这些能力允许对执行进行更精细的控制,并更好地与现有技术基础设施集成。
Execution IDs¶
ADK 框架根据父 ID 和计数器为子节点执行生成确定性标识符(ID)。ADK 工作流使用确定性 ID 来识别每个已调度节点的先前结果。这些 ID 根据动态节点调度的顺序生成,用于检查点以及在恢复或重新运行工作流时按正确顺序重新运行任务。
Custom execution IDs¶
在一些罕见的情况下,你可能需要稳定的标识符,例如在处理可重排序的列表时。通常你应该避免这样做,因为这会影响工作流任务重试和流程恢复。具体来说,这些 ID 用于检查节点状态并在节点已运行时跳过执行。如果你提供自定义 ID,请确保它们对于工作流重新运行是确定性的,并且在逻辑上对输入保持相同。
警告:自定义执行 ID
避免创建自定义执行 ID。由于执行 ID 用于确定节点的执行顺序,自定义执行 ID 可能会在系统尝试在你的工作流中重新运行这些节点时导致问题。
from google.adk import Context
from google.adk.workflow import node
from pydantic import BaseModel
from typing import Any
import asyncio
class Order(BaseModel):
order_id: str
cart_items: list[Product]
@node(rerun_on_resume=True)
async def process_all_orders(ctx: Context, node_input: Any):
orders = await get_orders()
process_tasks = []
for order in orders:
# 使用 run_id 提供自定义标识符。
# 自定义 run_id 必须包含至少一个非数字字符,
# 以避免与自动生成的顺序数字 ID 冲突。
task = ctx.run_node(process_order, order, run_id=f"order-{order.order_id}")
process_tasks.append(task)
results = await asyncio.gather(*process_tasks)
return results
默认情况下,自动生成的运行 ID 是从 "1" 开始的顺序整数(以字符串表示)。自定义 run_id 值必须包含至少一个非数字字符,以避免与这些自动生成的 ID 冲突。
将 runId 作为尾部选项传递给 ctx.runNode()。ID 必须包含至少一个非数字字符,
以避免与自动生成的顺序 ID 冲突:
import { node, NodeContext, Workflow } from '@google/adk';
interface Order {
orderId: string;
cartItems: string[];
}
/** Stands in for loading orders from a database. */
async function getOrders(): Promise<Order[]> {
return [
{ orderId: 'a91', cartItems: ['keyboard', 'mouse'] },
{ orderId: 'b02', cartItems: ['monitor'] },
{ orderId: 'c73', cartItems: ['dock', 'cable', 'hub'] },
];
}
const processOrder = node(
(_ctx: NodeContext, order: Order) =>
`order ${order.orderId}: ${order.cartItems.length} item(s) shipped`,
{ name: 'process_order' },
);
const processAllOrders = node(
async (ctx: NodeContext) => {
const orders = await getOrders();
const processTasks = orders.map((order) =>
ctx.runNode(processOrder, order, { runId: `order-${order.orderId}` }),
);
const results = await Promise.all(processTasks);
return results.map((result) => result.output).join('\n');
},
{ name: 'process_all_orders', rerunOnResume: true },
);
export const rootAgent = new Workflow({
name: 'root_agent',
edges: [['START', processAllOrders]],
});
在 Go 中,将 workflow.WithRunID("order-x") 作为尾部选项传递给 workflow.RunNode。ID 必须包含至少一个非数字字符,以避免与自动生成的顺序计数器 ID 冲突:
// newCustomIDWorkflow demonstrates supplying stable custom run IDs via
// workflow.WithRunID — equivalent to Python's:
//
// task = ctx.run_node(process_order, order, run_id=f"order-{order.order_id}")
//
// Custom run IDs must contain at least one non-numeric character to avoid
// collision with auto-generated sequential integer IDs.
func newCustomIDWorkflow() (agent.Agent, error) {
processOrderNode := workflow.NewFunctionNode("process_order",
func(_ agent.Context, orderID string) (string, error) {
return fmt.Sprintf("processed order %s", orderID), nil
},
workflow.NodeConfig{},
)
orders := []string{"ord-001", "ord-002", "ord-003"}
processAllOrders := workflow.NewDynamicNode[any, []string]("process_all_orders",
func(ctx agent.Context, _ any, _ func(*session.Event) error) ([]string, error) {
results := make([]string, 0, len(orders))
for _, orderID := range orders {
// WithRunID supplies a stable, deterministic identifier for
// each child invocation. IDs must contain at least one
// non-numeric character to avoid collision with the
// auto-generated sequential counter IDs.
result, err := workflow.RunNode[string](
ctx,
processOrderNode,
orderID,
workflow.WithRunID(fmt.Sprintf("order-%s", orderID)),
)
if err != nil {
return nil, fmt.Errorf("process order %s: %w", orderID, err)
}
results = append(results, result)
}
return results, nil
},
workflow.NodeConfig{},
)
return workflowagent.New(workflowagent.Config{
Name: "custom_id_workflow",
Description: "Processes orders with stable per-order execution IDs.",
Edges: workflow.Chain(workflow.Start, processAllOrders),
})
}