news 2026/9/13 19:00:41

ADK Python 工作流动态节点实战:用 ctx.run_node 实现运行时决定的执行路径

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
ADK Python 工作流动态节点实战:用 ctx.run_node 实现运行时决定的执行路径

ADK Python 工作流动态节点实战:用 ctx.run_node 实现运行时决定的执行路径

【免费下载链接】adk-pythonAn open-source, code-first Python toolkit for building, evaluating, and deploying sophisticated AI agents with flexibility and control.项目地址: https://gitcode.com/GitHub_Trending/ad/adk-python

本文围绕 ADK(Agent Development Kit,Python 版)的 Workflow 引擎讲解动态节点执行(Dynamic Node Execution)机制:当工作流的执行路径或节点运行次数无法在静态edges中确定时,如何用ctx.run_node配合原生 Python 控制流(如while循环)在运行时动态调度节点。基于contributing/samples/workflows/dynamic_nodes官方示例完整走读实现,并结合src/google/adk/workflowsrc/google/adk/agents/context.py的源码说明其底层调度、可恢复执行(rerun_on_resume)原理与事件轨迹验证方法。读完后你将掌握:如何声明动态编排节点、run_node各参数的实际影响、动态节点被中断后恢复的机制,以及如何用会话事件 JSON 验证多轮动态执行。

为什么需要动态节点执行

在 ADK Workflow 中,标准的执行路径由edges静态定义:哪些节点先跑、谁触发谁、走哪个分支,都在建图时写死。但真实场景中存在一类需求——具体的执行节点集合、或者某个节点要运行几次,只有在运行时才能确定,例如:

  • 循环生成内容直到通过 LLM 评审(评审轮数未知);
  • 根据用户输入的数量动态扇出子任务(fan-out 数量未知);
  • 动态拼接多个 Agent 的调用顺序。

dynamic_nodes示例处理的正是第一种:动态循环场景。一个名为orchestrate的 Python 节点充当驱动器(driver),在while True:循环中先执行generate_headlineAgent 基于给定主题生成头条,再执行evaluate_headlineAgent 对其评分。若评分为"tech-related"则返回该头条;若为"unrelated",则把评审反馈写回 state,循环继续。

值得注意的是,这个示例是标准loop示例的重写版本:它没有使用复杂的图边路由(例如在edges里配置条件路由函数),而是完全依靠原生 Python 控制流(while循环)加异步ctx.run_node调用来实现。执行拓扑如下:

图中只有START → orchestrate一条静态边,其余执行路径全部是orchestrate在运行时通过ctx.run_node动态调度出来的。

示例可运行的输入(见 README):

  • flower
  • quantum mechanics
  • renewable energy

完整示例代码解读

下面给出 agent.py 的完整实现,它是可独立运行的最小示例:

from typing import AsyncGenerator from typing import Literal from google.adk import Agent from google.adk import Context from google.adk import Event from google.adk import Workflow from google.adk.workflow import node from pydantic import BaseModel from pydantic import Field class Feedback(BaseModel): grade: Literal["tech-related", "unrelated"] = Field( description=( "Decide if the headline is related to technology or software" " engineering." ), ) feedback: str = Field( description=( "If the headline is unrelated to technology, provide feedback on how" " to make it more tech-focused." ), ) generate_headline = Agent( name="generate_headline", instruction=""" Write a headline about the topic "{topic}". If feedback is provided, take it into account. The feedback: {feedback?} """, ) evaluate_headline = Agent( name="evaluate_headline", instruction=""" Grade whether the headline is related to technology or software engineering. """, output_schema=Feedback, output_key="feedback", ) @node(rerun_on_resume=True) async def orchestrate( ctx: Context, node_input: str ) -> AsyncGenerator[Event | str, None]: yield Event(state={"topic": node_input}) while True: headline = await ctx.run_node(generate_headline) feedback = Feedback.model_validate( await ctx.run_node(evaluate_headline, node_input=headline) ) if feedback.grade == "tech-related": yield headline break root_agent = Workflow( name="root_agent", edges=[("START", orchestrate)], )

代码分为四个关键部分:

  1. 结构化评审输出Feedback:用 Pydantic 模型约束评审结果,grade字段取值限定为Literal["tech-related", "unrelated"],保证循环终止条件可被可靠判断。
  2. 生成 Agentgenerate_headline:instruction 中用{topic}从 state 读取主题,{feedback?}表示可选字段——首轮没有反馈时为空,后续轮次自动带上上一轮的改进建议。
  3. 评审 Agentevaluate_headlineoutput_schema=Feedback使其输出被解析为结构化对象;output_key="feedback"把结果写入会话 state 的feedback键,从而成为下一轮generate_headline可用的上下文。
  4. 编排节点orchestrate+ 根 Workfloworchestrate是唯一的静态节点,Workflow只声明了("START", orchestrate)一条边,其余全部动态发生。

关键机制一:@node(rerun_on_resume=True)启用可恢复执行

README 的第一步要点:要让 Python 节点使用ctx.run_node,必须声明为@node(rerun_on_resume=True)。它告诉引擎:如果任何被动态调度的子节点被中断(例如等待 human-in-the-loop 输入),工作流引擎应当暂停并在恢复时重新运行(re-run)编排节点,让编排节点从子节点拿到续跑后的结果。

这不是可选优化而是硬性要求。在源码 src/google/adk/workflow/_dynamic_node_executor.py 的run_node_internal入口处就有显式校验:

if not ctx._node_rerun_on_resume: raise ValueError( 'A node must have rerun_on_resume=True. Reason is that dynamically' ' scheduled nodes might be interrupted, and the workflow' ' wakes-up/re-runs the parent node, so it can get the child node' ' response.' )

即若父节点没有开启rerun_on_resume却调用ctx.run_node,会在运行期直接抛出ValueError。这个设计的原理是:动态子节点可能长时间挂起(HITL、长时工具),编排协程不能无限等待;引擎选择"暂停父节点 → 子节点完成/恢复 → 重跑父节点"的恢复模型,父节点重跑时从已完成子节点拿到缓存输出(去重),从而安全地继续循环。

@node装饰器本身还支持更多参数,完整签名见 src/google/adk/workflow/_node.py:name(覆盖节点名)、rerun_on_resumeretry_config(重试策略)、timeoutparallel_worker/max_parallel_workers(并行工作器)、auth_config(运行前请求用户认证,要求配合rerun_on_resume=True)、parameter_binding'state'默认从ctx.state绑定函数参数,'node_input'则从node_input绑定并按函数签名推断输入/输出 schema,用于节点充当 Agent 工具的场景)。动态节点场景中,与本例直接相关的是rerun_on_resume

关键机制二:ctx.run_node从上下文运行节点

README 的第二步要点:向 Python 节点注入ctx: Contextawait ctx.run_node(node_to_run)运行目标节点;返回值就是该次执行的最终输出;在循环的每次迭代之间还可以yield Event(...)更新 state。示例中:

  • await ctx.run_node(generate_headline)—— 不传node_input,Agent 直接从 state 读取{topic}/{feedback?}生成头条,返回值即头条文本;
  • await ctx.run_node(evaluate_headline, node_input=headline)—— 把本轮头条作为node_input显式传入评审 Agent,返回值经Feedback.model_validate解析为结构化对象;
  • yield Event(state={"topic": node_input})—— 在第一轮循环前把用户输入写入 state,让下游 Agent 的 instruction 模板能取到topic

Context.run_node的公开 API 定义在 src/google/adk/agents/context.py,完整参数与语义如下(均可用于动态节点编排):

参数说明
node要执行的节点:BaseNode实例,或任何可被构建成节点的可调用对象/Agent
node_input传给该节点的输入数据,默认None
use_as_output若为True,子节点输出直接作为调用节点的输出,调用节点自身的输出事件被抑制以避免重复
run_id自定义本次动态执行的 run ID,便于跨运行关联事件;不提供则自动生成
use_sub_branch若为True,子节点在子分支中执行,隔离其 state 与事件
override_branch覆盖父节点默认分支
override_isolation_scope覆盖父节点默认的隔离作用域
raise_on_wait若为True,子节点处于 WAITING 时抛出NodeInterruptedError而非返回None

该方法文档中有一条重要的使用约束:必须直接await,不要用asyncio.create_task()包裹——那样任务将失去监督:错误会被静默吞掉,父节点被中断(如 HITL)时任务也不会被取消。本示例中两次调用都是直接await,符合这一要求。

底层调度:动态节点如何被引擎接管

从源码结构看,ctx.run_node最终进入 run_node_internal,其核心逻辑分三种路径:

  1. 前置校验(见上文):强制父节点rerun_on_resume=True
  2. 双模式执行
    • Workflow 模式:当父节点运行在工作流图内(ctx._workflow_scheduler存在),执行委托给工作流调度器(ScheduleDynamicNode),由调度器处理图依赖与状态。这里还有一个细节校验——调用方显式传入的run_id如果是纯数字会被拒绝("must contain non-numeric characters to prevent collision with auto-generated IDs"),因为数字 run_id 由调度器顺序分配,显式数字会与之冲突。
    • Standalone 模式:节点独立于工作流运行时,直接通过NodeRunner执行。
  3. Agent Transfer 循环run_node_internal用一个while True循环处理"节点执行中又请求转移到另一个 Agent"的情况(例如子 Agent A 把执行转给 Agent B),循环内更新目标节点与父上下文指针并继续;若子节点报错则抛DynamicNodeFailError,若子节点被中断(存在interrupt_ids)则把中断 ID 传播回父节点上下文并抛NodeInterruptedError,让上层 NodeRunner 把父节点记录为等待状态而非误判完成——这正是rerun_on_resume恢复模型能工作的前提。

也就是说,示例中"看起来只是两行await ctx.run_node(...)"的背后,引擎完成了 run_id 分配、分支/作用域隔离、事件溯源(子节点输出事件带上output_for指向父节点路径)、中断传播与恢复重放这一整套机制。

用事件轨迹验证多轮动态执行

示例目录附带了一个真实的会话事件文件 tests/flower.json,记录了输入flower时的一次完整运行,可以把它当作"动态节点如何落事件"的参照:

  • e-2orchestrate产生stateDelta: {"topic": "flower"},路径为root_agent@1/orchestrate@1
  • e-3:第一次动态调度generate_headline,输出"A World of Petals",其nodeInfo.pathroot_agent@1/orchestrate@1/generate_headline@1,且output_for指回父节点路径——即"子节点输出归属于父节点"的体现;
  • e-4:第一次evaluate_headline输出结构化 JSON,gradeunrelatedstateDelta写入feedback(含改进建议);
  • e-5/e-6:第二轮循环,generate_headline的路径变为.../generate_headline@2(run 序号递增),结合反馈产出"**Digital Petals: Engineering AI-Enhanced Blooms**"evaluate_headline@2评出tech-related
  • e-7:最终输出即该头条,output_for同时指向orchestrate与根root_agent

这段轨迹印证了三个机制:state 通过Event(state=...)在轮次间传递(feedback每轮被覆写);动态子节点的 path/run_id 由引擎自动生成并随轮次递增;循环的终止(yield headlinebreak)最终把父节点输出冒泡为整个 Workflow 的输出。

测试覆盖与延伸阅读

仓库的单元测试 tests/unittests/workflow/test_workflow_dynamic_nodes.py 系统地覆盖了动态节点调度的三种核心情形——全新执行(无历史事件)、已完成去重(重跑时直接返回缓存输出)、中断后恢复(携带resume_inputs重跑),以及多动态节点、嵌套动态节点、use_as_output输出委托等边界情况。此外还有 test_dynamic_node_executor.py、test_dynamic_node_scheduler.py、test_dynamic_use_as_output.py 等针对执行器与调度器的独立测试,可作为理解上述底层机制的对照材料。

实践要点小结

  • 静态边最少化:动态节点场景下Workflow只需声明("START", 编排节点)一条边,其余拓扑由 Python 控制流表达,可读性和灵活性优于在edges中写条件路由函数;
  • 必须rerun_on_resume=True:这是调用ctx.run_node的硬性前提,用于支撑"子节点中断 → 父节点重跑"的恢复模型,源码中有显式校验;
  • 直接await ctx.run_node:不要用asyncio.create_task包裹,否则错误与中断传播都会失效;
  • 用 state 做轮次间的通信总线:通过yield Event(state=...)写入topic、利用output_key写回feedback,让多轮循环中的 Agent 自然获得上一轮上下文;
  • output_schema固化循环终止条件:结构化的grade字段让while循环的退出判断可靠、可测试;
  • 用事件 JSON 验证行为:动态子节点的 path(@序号递增)与output_for归属关系是核对多轮执行是否正确落地的直接证据。

适用前提与限制:本机制面向 ADK 的 Workflow 引擎(google.adk.Workflow/@node),需要可被build_node构建的节点对象(BaseAgent、BaseTool、BaseNode 或可调用对象);显式传入纯数字run_id会被拒绝;use_as_output委托在非 Workflow 父节点下每个父节点只能设置一次。以上行为均以当前仓库源码为准。

【免费下载链接】adk-pythonAn open-source, code-first Python toolkit for building, evaluating, and deploying sophisticated AI agents with flexibility and control.项目地址: https://gitcode.com/GitHub_Trending/ad/adk-python

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/13 18:57:35

通用MCU+硅MOS做FOC驱动的硬件瓶颈深度解析

1. 项目概述:为什么“通用MCU 硅MOS”在FOC驱动中总卡在体积与扭矩的矛盾点上?你有没有拆过市面上那些标称“300W无刷电机驱动板”,尺寸比名片还小,却能带动2kgcm以上堵转扭矩的负载?我去年帮一家电动工具客户做竞品逆…

作者头像 李华
网站建设 2026/9/13 18:57:26

MSO算法在无人机路径规划中的Matlab实现与应用

1. 项目概述:MSO算法与无人机路径规划2025年算法海市蜃楼算法(Mirage Simulation Optimization,简称MSO)是新一代基于环境动态模拟的智能路径规划方法。这个算法最有趣的特点在于它能模拟出类似"海市蜃楼"的虚拟环境扰动…

作者头像 李华
网站建设 2026/9/13 18:52:57

相关杂波生成与ZMNL方法:雷达海杂波仿真的关键

简介:面向无线通信与雷达系统中的相关杂波建模,MATLAB仿真资源包聚焦多类统计模型,适用于信号处理、通信工程等领域的研究生与研发工程师,可用于生成和分析多种统计分布的杂波场景。压缩包内共9个m文件,均为可直接运行…

作者头像 李华