动态智能体工作流¶
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 动态工作流的上下文中运行。
在 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,
)
手动创建节点包装器代码在以下情况会很有用:当你要包装来自外部库的函数时,需要从同一函数创建具有不同配置的多个节点时,或者当你要在注册表中管理节点引用以进行高级编排时。
在 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)],
)
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
在 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 动态工作流创建顺序任务处理。
以下代码片段展示了一个动态工作流,包含一个智能体、一个函数节点和第二个智能体:
在 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.",
output_schema=str,
)
@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}""",
output_schema=str,
)
@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)
return code
在 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
提示:恢复并行节点
工作流框架确保如果动态工作流被恢复,只有失败或中断的工作节点会被重新执行,包括并行工作节点。
在 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 以正确处理中断。
在 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 冲突。
在 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),
})
}