Agno Function Workflow 实战指南:用单个可调用函数编排多智能体工作流
【免费下载链接】agnoBuild, run, and manage agent platforms.项目地址: https://gitcode.com/GitHub_Trending/ag/agno
本指南以 cookbook/04_workflows/01_basic_workflows/03_function_workflows 目录下的演示为核心,讲解 Agno(agno)工作流(Workflow)中一种独特的编排形态——函数工作流(Function Workflow):不再把步骤声明成显式的Step对象列表,而是将一个可调用函数整体赋给Workflow(steps=...),让编排逻辑完全收拢在普通 Python 函数里。读完本文,你将掌握函数工作流的函数签名约束、WorkflowExecutionInput数据结构,以及同步 / 流式 / 异步 / 异步流式四种运行模式的写法,并理解其在Workflow源码层的实现机制。
一、什么是 Function Workflow
在 Agno 的 cookbook 工作流体系中,01_basic_workflows 目录之下并行组织了三种基础的流程形态:
01_sequence_of_steps/:以Step/Steps对象列表串联的顺序执行工作流;02_step_with_function/:在单个步骤内使用函数类型执行器;03_function_workflows/:整体用一个执行函数替代步骤列表,即本节主题。
该目录的 README.md 对范围的定义非常简洁——"Runnable workflow examples under 03_function_workflows",并指明唯一示例文件function_workflow.py的作用是Demonstrates function workflow。换句话说,这个示例的实质是:Workflow的steps参数除了接受Step列表、Steps容器之外,还接受单个可调用对象(callable);当传入可调用对象时,每次运行工作流都会调用这个函数,由函数内部自行完成对 Agent、Team 乃至其他 Workflow 的编排。
在 Workflow 实现 中可以看到steps字段的声明(约第 623 行),而在运行参数解析、序列化等多个逻辑分支里都专门针对callable(self.steps)做了区分处理(如run_parameters属性、状态字典序列化等),可以确认"可调用 steps"是 Workflow 一等公民的运行方式。
二、运行环境与前置条件
README.md 明确给出了两个前置条件,这也是整个 cookbook 的运行约定:
- 激活 demo 虚拟环境:仓库统一使用
.venvs/demo/bin/python运行示例(该环境由仓库根目录 scripts/demo_setup.sh 之类的脚本创建); - 加载 API Key:使用
direnv allow放行本地.envrc文件,为示例中的模型与联网工具注入密钥。
运行示例:
# 在仓库根目录执行 .venvs/demo/bin/python cookbook/04_workflows/01_basic_workflows/03_function_workflows/function_workflow.py脚本在if __name__ == "__main__":块中会依次跑完四种模式(同步、同步流式、异步、异步流式),因此一次执行即可观察到函数工作流在不同运行模式下的完整行为。
三、示例全景:一个内容策划流水线
function_workflow.py 构造了一个"研究 → 策划"两阶段内容生产流水线,其组件布局如下:
| 组件 | 类型 | 职责 |
|---|---|---|
hackernews_agent | Agent | 使用HackerNewsTools(),从 Hacker News 帖子中提炼洞察 |
web_agent | Agent | 使用WebSearchTools(),检索最新资讯与趋势 |
research_team | Team | 成员为上述两个 Agent,统一做技术主题调研 |
content_planner | Agent | 依据调研结果生成 4 周内容排期 |
streaming_hackernews_agent | Agent | 流式模式专用的调研 Agent |
组装方式如下(模型均使用OpenAIChat,示例中取id="gpt-5.6-luna"/"gpt-5.2",读者按自己 API Key 可用模型替换即可):
research_team = Team( name="Research Team", members=[hackernews_agent, web_agent], instructions="Research tech topics from Hackernews and the web", )四个Workflow实例共享同一套数据库与步骤函数,但各自绑定一种模式的执行函数:
sync_workflow = Workflow( name="Content Creation Workflow", description="Automated content creation from blog posts to social media", db=SqliteDb(session_table="workflow_session", db_file="tmp/workflow.db"), steps=custom_execution_function, )这里值得注意两点:
steps=直接接收函数引用而非步骤列表,这正是"函数工作流"的命名由来;db=SqliteDb(session_table=..., db_file=...)让工作流运行数据持久化到本地 SQLite,为函数工作流提供会话与断点恢复能力。
四、执行函数签名与 WorkflowExecutionInput
函数工作流的执行函数遵循统一签名约定。以示例中的同步版本为例:
def custom_execution_function( workflow: Workflow, execution_input: WorkflowExecutionInput, ) -> str: print(f"Executing workflow: {workflow.name}") run_response = research_team.run(execution_input.input) research_content = run_response.content planning_prompt = f"""...Core Topic: {execution_input.input} Research Results: {research_content[:500]}...""" content_plan = content_planner.run(planning_prompt) return content_plan.content4.1 两个固定参数
函数前两个位置参数由 Workflow 自动注入:
workflow:当前Workflow实例,可读取name等配置;execution_input:WorkflowExecutionInput类型,封装本次运行的全部输入。
从源码层的run_parameters属性可以看到,框架在解析执行函数签名时,会显式剔除workflow、execution_input(以及self)这三个参数名,把它们当作框架保留参数;其余参数则被识别为工作流的运行参数(见 workflow.py 第 866-884 行附近的signature(self.steps)内省逻辑)。这意味着你可以在执行函数里声明额外参数,由调用方在运行工作流时按名传入,实现参数化。
4.2 WorkflowExecutionInput 的字段
WorkflowExecutionInput定义于 agno/workflow/types.py(数据类约从第 248 行开始),核心字段包括:
| 字段 | 类型 | 说明 |
|---|---|---|
input | str/Dict/List/BaseModel | 工作流主输入,示例中即"AI trends in 2024" |
additional_data | Dict[str, Any] | 附加上下文数据 |
images/videos/audio/files | 各媒体类型列表 | 多模态输入载体 |
它还提供了get_input_as_string()方法,能将字符串、Pydantic 模型、字典或列表统一序列化为字符串,方便你在函数内部把输入拼进 Prompt。对比同一文件中的StepInput(面向显式步骤,携带previous_step_outputs、workflow_session等跨步骤数据)可以看出:函数工作流的执行函数是"黑盒式"的——它只接收本次运行的输入,步骤间的数据流转完全由函数内部自行组织,这正是该形态灵活性的来源。
五、四种运行模式:sync / stream / async / async-stream
示例的进阶价值在于用四个形态相近的执行函数覆盖了 Agno 支持的全部运行模式。逐一拆解:
5.1 同步(sync)
返回值str,内部依次调用research_team.run()与content_planner.run():
def custom_execution_function(workflow, execution_input) -> str: ...工作流侧用sync_workflow.print_response(input="AI trends in 2024")触发。
5.2 同步流式(sync streaming)
返回类型标注为Iterator,关键差异在内部:
def custom_execution_function_stream(workflow, execution_input) -> Iterator: research_content = "" for response in streaming_hackernews_agent.run( execution_input.input, stream=True, stream_events=True, ): if hasattr(response, "content") and response.content: research_content += str(response.content) ... yield from content_planner.run(planning_prompt, stream=True, stream_events=True)要点有二:
- 上游 Agent 开启
stream=True, stream_events=True逐条消费事件,把content累加为完整调研文本后再构造下游 Prompt(这也是它单独准备streaming_hackernews_agent的原因); - 末尾用
yield from把content_planner.run(..., stream=True)的事件流原样透传给调用方。
工作流侧通过sync_stream_workflow.print_response(input=..., stream=True)触发逐块输出。
5.3 异步(async)
函数以async def定义,返回str,内部使用await content_planner.arun(planning_prompt):
async def custom_execution_function_async(workflow, execution_input) -> str: ... content_plan = await content_planner.arun(planning_prompt) return content_plan.content运行侧需要用事件循环驱动:
asyncio.run(async_workflow.aprint_response(input="AI trends in 2024"))5.4 异步流式(async streaming)
返回类型为AsyncIterator,上游用async for消费arun(..., stream=True, stream_events=True),下游同样async for逐条yield:
async def custom_execution_function_async_stream(workflow, execution_input) -> AsyncIterator: async for response in streaming_hackernews_agent.arun(...): if hasattr(response, "content") and response.content: research_content += str(response.content) ... async for response in content_planner.arun(planning_prompt, stream=True, stream_events=True): yield response运行侧:
asyncio.run(async_stream_workflow.aprint_response(input="AI trends in 2024", stream=True))5.5 模式对照速查
| 模式 | 函数形态 | 返回类型 | 内部调用 | 工作流侧触发 |
|---|---|---|---|---|
| Sync | def | str | .run() | print_response() |
| Sync Streaming | def | Iterator | .run(stream=True)+yield from | print_response(stream=True) |
| Async | async def | str | await .arun() | asyncio.run(aprint_response()) |
| Async Streaming | async def | AsyncIterator | async for+yield | asyncio.run(aprint_response(stream=True)) |
规律非常清晰:同步用run/print_response,异步加a前缀变成arun/aprint_response;要流式就传stream=True并把返回类型切换为对应的Iterator/AsyncIterator。示例中四个 Workflow 均命名为 "Content Creation Workflow" 且复用同一 SQLite 会话表,进一步演示了同一业务在多运行形态下的平行复用。
六、执行流程剖析:一次运行内发生了什么
把custom_execution_function的代码与调用链对齐,一次同步运行的完整流程是:
- 用户在终端调用
sync_workflow.print_response(input="AI trends in 2024"); - Workflow 构造
WorkflowExecutionInput(input="AI trends in 2024")并调用执行函数(steps指向的可调用对象); - 执行函数把
research_team(HackerNews Agent + Web Agent)当作"第一步",run()返回的content作为中间产物; - 代码对
research_content做[:500]截断后嵌入精心构造的planning_prompt——这是控制上下文窗口、约束下游输入规模的实用技巧; - 执行函数再驱动
content_planner产出最终内容排期,并return content_plan.content; - Workflow 把返回值包装为运行结果,同时借助配置的
SqliteDb将会话写入tmp/workflow.db。
在流式变体中,第 3、5 步换为事件流逐条累计/透传,中间产物仍是纯字符串拼接,整体编排语义保持一致。可见函数工作流把"步骤"这一概念彻底函数化:没有显式的Step生命周期、没有步骤间的输入传递管线,一切顺序、分支、截断、Prompt 拼接都由你写在函数体里,获得最大的表达自由度。
七、函数工作流的适用场景与边界
结合源码结构可以做如下归纳(属于从代码与示例得出的推断性结论):
适合函数工作流的场景:
- 编排逻辑高度定制、不便拆成标准
Step的流程; - 希望把"调研团队 + 规划 Agent"这类已有 Agent / Team 组合直接封装为可复用流水线的场景;
- 需要以统一函数接口同时对外提供 sync / async / streaming 多形态 API 的服务层。
需要注意的边界:
- 函数工作流内部对 Agent / Team 的调度(如并行、重试、循环)需要自行在函数内实现,它不享受显式步骤容器的调度便利(后者对应
Steps/Parallel/Loop/Condition等组件); - 如果业务需要基于步骤粒度做会话恢复、指标聚合与人工审批(HITL),更适合回到显式步骤路线;函数工作流的持久化粒度在"整个函数"这一层。
八、进一步探索
- 示例源码:function_workflow.py(完整对照阅读四种执行函数写法);
- 同目录说明:03_function_workflows/README.md;
- 相邻形态:01_sequence_of_steps(显式步骤序列)、02_step_with_function(步骤内函数)可用于对比理解三种流程定义风格的差异;
- 核心数据结构:agno/workflow/types.py(
WorkflowExecutionInput等定义); - 底层调度实现:agno/workflow/workflow.py(
steps参数对可调用对象的解析、签名内省与运行分发)。
函数工作流是 Agno 工作流体系中"以代码为编排核心"的代表性写法。它牺牲了显式步骤带来的结构化能力,换来的是完全自由的函数体编排与极低的抽象负担——当你手中的 Agent、Team 已经足够"聪明",而流程本身又不足以复杂到需要一张步骤图时,它就是最贴合直觉的选择。
【免费下载链接】agnoBuild, run, and manage agent platforms.项目地址: https://gitcode.com/GitHub_Trending/ag/agno
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考