Agno Workflow 后台执行实战:异步轮询与 WebSocket 实时事件流
【免费下载链接】agnoBuild, run, and manage agent platforms.项目地址: https://gitcode.com/GitHub_Trending/ag/agno
导读
本篇技术指南以 cookbook/04_workflows/06_advanced_concepts/background_execution 目录为骨架,系统讲解如何在 Agno 中把 Workflow 放到后台运行(background=True)并通过两种方式取回结果:一是非流式模式下基于run_id的轮询(poll),二是流式模式下通过 WebSocket/SSE 实时推送RunStarted、StepStarted、RunContent等运行事件。读完你不仅能写出异步后台工作流,还能搭建一套「FastAPI + WebSocket 服务端 + Rich 交互式客户端」的完整后台执行示例。
一、示例目录定位与总体结构
该目录隶属于「04_workflows/06_advanced_concepts」高级概念系列,在父级说明中(见 cookbook/04_workflows/06_advanced_concepts/README.md)被定位为background_execution:后台执行补充示例,与long_running(长任务)、run_control(运行控制)等主题并列,面向已经掌握 Workflow 基础(定义 Step、串联 Agent/Team)的读者。
目录下共有三个可直接运行(runnable)的 Python 文件与两个文档文件,README 给出的职责划分如下:
| 文件 | 演示内容 |
|---|---|
| background_poll.py | 异步后台运行 Workflow,并轮询运行状态直到完成 |
| websocket_client.py | 演示 WebSocket 客户端(认证、启动工作流、渲染流式事件) |
| websocket_server.py | 演示 WebSocket 服务端(后台运行工作流并把事件推给客户端) |
三者恰好覆盖了后台执行的完整闭环:轮询方案负责「起一个后台任务 → 主动查询」,WebSocket 方案负责「服务端推送 → 客户端被动接收并展示」。
运行前的前置条件(README 原文)包括:
- 激活演示环境:
.venvs/demo/bin/python; - 用
direnv allow加载 API 密钥(需要本地.envrc文件); - 部分示例依赖本地 AgentOS 服务,具体服务地址见示例文件头部或运行打印(例如 WebSocket 服务地址为
ws://localhost:8000/ws)。
二、原理先行:Workflow.arun的background参数与三种后台模式
三个示例都建立在同一个入口arun()之上。从 libs/agno/agno/workflow/workflow.py#L10793-L10922 的签名可以看到,arun在常规入参之外额外暴露了四个与后台执行直接相关的参数:
| 参数 | 类型 | 作用 |
|---|---|---|
background | Optional[bool] = False | 是否以后台方式启动运行 |
stream | Optional[bool] = None | 是否流式返回内容 |
stream_events | Optional[bool] = None | 是否同时推送运行事件(RunStarted/StepStarted 等) |
websocket | Optional[WebSocket] | 显式传入 WebSocket 连接;同时会开启enable_websocket |
enable_websocket | bool = False | 与background+stream同时为 True 时,改用 WebSocket 传输代替默认 SSE 传输 |
源码中background=True时会按以下优先级路由到三种底层实现:
background=True+stream=True+ WebSocket 启用(传了websocket或enable_websocket=True)→_arun_background_stream_ws:实时事件 + WebSocket 传输。代码中注释明确写着 "Background + Streaming + WebSocket = Real-time events (opt-in)";background=True+stream=True(未启用 WebSocket)→_arun_background_stream:后台 + 默认 SSE 传输;background=True+ 非流式 →_arun_background:后台 + 轮询模式,即 background_poll.py 走的分支。
另外注意源码做了向后兼容处理:只要传入了websocket参数(即使没写enable_websocket),也会自动把enable_websocket置为True(见 workflow.py#L10835-L10837)。stream_events在两种流式分支下都会被开启或显式保留,因此客户端能收到结构化的运行生命周期事件。
需要区分的是:run()与arun()是同步/异步两套重载;后台执行全部走arun(异步实现),轮询读结果则用get_run()/aget_run()。
三、方案一:后台运行 + 轮询(background_poll.py)
background_poll.py 的场景非常典型:研究 Hacker News 与 Web 上的科技话题,再由内容策划 Agent 输出为期四周的内容排期。整个过程耗时较长,因此示例选择「异步发起 + 定时轮询」而不是同步阻塞等待。
3.1 组装 Agent、Team 与 Step
示例先创建三个 Agent:Hackernews Agent(工具HackerNewsTools)、Web Agent(工具WebSearchTools)负责研究,Content Planner按指令规划内容排期;再把两个研究 Agent 放进一个Team(name="Research Team")。随后把两个执行单元各包成一个 Step:
research_step = Step(name="Research Step", team=research_team) content_planning_step = Step(name="Content Planning Step", agent=content_planner)Step与Workflow的类型定义位于agno.workflow.step与agno.workflow.workflow。这里的核心思路是:Agent 是单一能力执行者,Team 是横向协作单元,Step 是 Workflow 的最小编排节点——Step 既可挂 Agent 也可挂 Team,编排层只需关心 Step 的先后顺序。
3.2 用 SqliteDb 持久化会话
Workflow 构造时传入了一个SqliteDb(SQLite 会话数据库,演示环境的库也覆盖在 cookbook/06_storage/sqlite 等示例中):
content_creation_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=[research_step, content_planning_step], )配置项含义:
session_table:会话表名,这里用workflow_session;db_file:SQLite 文件落盘路径,tmp/workflow.db会在运行时创建。
为什么后台模式必须配 db?后台任务发起后,当前协程无法直接拿到最终结果,只能通过run_id反查;而反查的数据源正是会话数据库。没有 db,就无法跨请求/跨进程恢复运行状态。
3.3 发起后台任务并立即拿到 Initial Response
bg_response = await content_creation_workflow.arun( input="AI trends in 2024", background=True, ) print(f"Initial Response: {bg_response.status} - {bg_response.content}") print(f"Run ID: {bg_response.run_id}")background=True时arun立即返回一个「初始响应」,其中值得关注三个字段:
run_id:本次运行的唯一标识,轮询阶段的查询键;status:当前运行状态(如运行中/排队中),此时打印的并非最终结果;content:由于是后台模式,这里通常还是空/占位内容——不能把 initial response 当作执行结果。
3.4 轮询循环:get_run + has_completed + 超时保护
发起后主流程进入循环,每 5 秒轮询一次:
while True: poll_count += 1 print(f"\nPoll #{poll_count} (every 5s)") result = content_creation_workflow.get_run(bg_response.run_id) if result is None: print("Workflow not found yet, still waiting...") if poll_count > 50: # 尚未入库的重试上限 print(f"Timeout after {poll_count} attempts") break await asyncio.sleep(5) continue if result.has_completed(): # 运行完成的判据 break if poll_count > 200: # 总轮询上限(约 1000 秒) print(f"Timeout after {poll_count} attempts") break await asyncio.sleep(5)示例设计了两级超时保护,值得写生产代码时借鉴:
result is None表示run_id还没在库中落账(后台任务可能仍在初始化),最多重试 50 次;result.has_completed()为 False 且轮询超过 200 次则强制退出,避免无限等待。
get_run()的语义可以从实现确认(见 workflow.py#L5659 附近):它按run_id(可叠加session_id)从会话数据库读取运行记录并返回WorkflowRunOutput,注释明确标注这是获取后台运行状态与细节的简化接口。实现里还有一个重要约束——同步数据库用get_run(),异步数据库必须改用aget_run(),否则会抛出ValueError。
3.5 输出最终结果
循环结束后再查一次拿到完整结果,并用 Agno 提供的pprint_run_response(result, markdown=True)美化打印——它会把 run_id、会话信息与各步最终内容以 Markdown 形式渲染出来。
整个文件以asyncio.run(main())驱动,说明这套轮询逻辑天然适配 FastAPI/AgentOS 等异步运行环境。
四、方案二:服务端——后台工作流事件经 WebSocket 推送(websocket_server.py)
websocket_server.py 基于FastAPI + uvicorn,把「后台运行 Workflow + 事件流式推送」封装成一个可被任意客户端连接的服务。
4.1 服务拓扑与启动
- 启动后监听
0.0.0.0:8000; - WebSocket 端点:
ws://localhost:8000/ws; - HTTP 状态端点:
GET /,返回status、endpoints、当前连接数connections与已认证数authenticated; - 附带 FastAPI 原生文档:
http://localhost:8000/docs。
服务端维护两个全局字典:
active_connections: Dict[str, WebSocket] = {} authenticated_connections: Dict[str, bool] = {} # {connection_id: is_authenticated}每个连接获得一个自增的connection_id(形如conn_0),建立时默认未认证,断线时在finally中清理两条记录。
4.2 基于 SECURITY_KEY 的握手认证
认证协议非常简单:客户端发送{"action": "authenticate", "token": "..."},服务端校验通过后回authenticated事件并标记该连接已认证;token 缺失回auth_error("Token is required"),错误 token 回auth_error("Invalid token")。未认证连接发送其它指令时会收到auth_required事件。
SECURITY_KEY = os.getenv("SECURITY_KEY", "your-secret-key") def validate_token(token: str) -> bool: if not SECURITY_KEY or SECURITY_KEY == "your-secret-key": return True # 未配置密钥时默认放行(演示环境行为) return token == SECURITY_KEY可见安全策略是「没设密钥就全放行,设了密钥则严格比对」——生产部署务必通过环境变量SECURITY_KEY配置真实密钥。
4.3 后台 + WebSocket 的调用要点
收到start-workflow消息后,服务端先为本次请求新建一个独立的 Workflow(两个 Step 分别挂研究 Agent 与搜索 Agent,会话持久化在tmp/workflow_bg.db、表名workflow_bg),然后这才是整个示例的精华:
result = await workflow.arun( input=message, session_id=session_id, stream=True, stream_events=True, background=True, websocket=websocket, )对照第二节的路由逻辑,这一调用组合(background + stream + websocket)会命中_arun_background_stream_ws,即后台执行、同时把运行事件实时写回 WebSocket。事件推送不是手写的——Agno 会构造一个WebSocketHandler包装传入的 FastAPIWebSocket(见 workflow.py#L10829-L10837),把后台运行过程中的事件序列化后逐个发给客户端。调用前后服务端再补发两类生命周期事件:
- 调用前:
workflow_starting(携带原始 message 与 session_id); - 成功后:
workflow_initiated(携带run_id、session_id,表示后台流式工作流已成功启动); - 异常时:
workflow_error。
这样客户端既能收到工作流的「控制面事件」(启动/完成/错误),又能收到 Agno 内核产生的「运行面事件」(各 Step 开始/结束、Token 级内容流、工具调用前后等)。
4.4 其余协议细节
ping→ 回pong;- 其它未识别消息 → 回
echo(便于联调); - 单条消息处理异常会回
error事件并把异常文本带给客户端,连接本身不关闭。
五、方案三:客户端——Rich 交互终端与事件渲染(websocket_client.py)
websocket_client.py 是一个基于websockets与rich的交互式客户端,用于连接方案二的服务端,其职责是:连接 → 认证 → 启动工作流 → 持续渲染服务端推送的事件。
5.1 命令行入口
# 交互模式(默认,无 message 时也进入交互模式) .venvs/demo/bin/python websocket_client.py -i # 单发模式:连上后立即用一句话启动工作流 .venvs/demo/bin/python websocket_client.py -m "AI trends 2024" # 自定义服务地址与认证 token(也支持 SECURITY_KEY 环境变量) .venvs/demo/bin/python websocket_client.py --server ws://localhost:8000/ws --token xxx -i参数一览:--server(默认ws://localhost:8000/ws)、--message/-m、--interactive/-i、--token/-t。token 的取值顺序是「命令行参数优先,否则读SECURITY_KEY环境变量」。
5.2 交互指令集
| 指令 | 行为 |
|---|---|
auth | 提示输入 token 并发起认证 |
start <message> | 用消息启动工作流,自动生成cli-session-<时间戳>会话 |
ping | 发送心跳并观察pong |
quit/exit/q | 退出并清理监听任务、断开连接 |
未认证时连接提示栏会高亮显示 "AUTHENTICATION REQUIRED",提醒先输入auth。
5.3 事件协议解析:JSON 与 SSE 双格式兼容
listen_for_events的解析策略体现了对两种服务端实现风格的兼容:
- 先尝试
json.loads整体解析(纯 JSON 事件,如方案二服务端发送的消息); - 解析失败则按SSE 文本格式二次解析:形如
event: X+data: {...}的多行消息,parse_sse_message会提取event_type并json.loads出 data,最后把type字段合并进 dict(这对应 Agno 默认 SSE 传输分支的输出)。
5.4 事件渲染与流式内容累积
客户端内置一张事件→样式的映射表,覆盖connected/authenticated/auth_error/WorkflowStarted/StepStarted/StepCompleted/WorkflowCompleted/WorkflowError/RunStarted/RunContent/RunCompleted/ToolCallStarted/ToolCallCompleted等十余种事件,每种都映射到标签与 Rich 颜色。
最有价值的是对RunContent 流式内容的累积渲染:current_step_content按step_id持续拼接每次到达的内容分片,且遵循「分片长度 > 3 或含换行才渲染」的节流策略,避免把单个字符刷成满屏面板;当累积文本超过 300 字符时,面板只显示最后 300 字符并加...前缀,防止终端被刷爆。同时每个事件面板还会附带step_name、agent_name、run_id、session_id、step_index等关键字段,帮助观察「哪一步、哪个 Agent 正在输出什么」。
这种设计非常适合把多 Step 工作流的实时进度做成 Web 终端或运维大屏。
六、串联运行与测试验证
6.1 推荐运行顺序
由于三者之间存在依赖关系,建议按「服务端 → 客户端 → 轮询」顺序联调:
# 终端 1:启动 WebSocket 服务端 .venvs/demo/bin/python cookbook/04_workflows/06_advanced_concepts/background_execution/websocket_server.py # 终端 2:以单发模式消费一个后台工作流 .venvs/demo/bin/python cookbook/04_workflows/06_advanced_concepts/background_execution/websocket_client.py \ --server ws://localhost:8000/ws -m "AI trends 2024"若选择运行 background_poll.py,该脚本不依赖 WebSocket 服务,只需在事件循环中自行完成后台发起 + 轮询;tmp/workflow.db会随运行自动创建。
6.2 测试日志给出的可执行性参考
目录下的 TEST_LOG.md 记录了仓库自动化测试时的三条实证结果,可作为运行预期:
background_poll.py:normal 模式执行 35 秒后超时(日志显示已完成多轮 Agent Run),说明完整跑完需要真实模型调用与更长执行窗口,评估超时阈值时应放宽;websocket_client.py:startup 模式通过(8 秒内完成启动校验),因当时无服务端而报连接失败Connect call failed,属预期行为——它必须在服务端存活时才能完整演示;websocket_server.py:startup 模式通过,8 秒后按预期终止进程,说明服务可正常拉起。
七、把方案落地到生产的关键清单
综合源码实现与示例细节,在真实项目中落地「后台执行」时可沉淀如下经验:
- 选对后台模式:只需要最终结果 → 非流式
background=True+get_run()轮询;需要实时进度 →background=True + stream=True + stream_events=True,并决定用默认 SSE 还是传入websocket启用 WebSocket 通道; - 务必配置会话数据库:
run_id反查依赖持久化(示例均用SqliteDb,可用 cookbook/06_storage/sqlite 中的 SQLite 系列作参照);注意异步数据库要改用aget_run(); - 把超时与重试写进轮询:参考 background_poll.py 的「未入库重试上限 + 总轮询上限」两级护栏;
- 安全边界:仿照服务端用
SECURITY_KEY做 token 鉴权,未认证连接只允许authenticate,其余指令一律回auth_required; - 客户端体验:用事件类型→样式的映射统一渲染,对
RunContent做按 step 累积与长度截断,避免大量小分片刷屏。
通过 background_execution 这一组示例,你掌握的不仅是两个 API 的用法,而是 Agno 后台执行从「发起 → 持久化 → 查询/推送 → 渲染」的完整链路——这正是把长耗时多 Agent 工作流接入 Web 服务与实时前端的标准姿势。
【免费下载链接】agnoBuild, run, and manage agent platforms.项目地址: https://gitcode.com/GitHub_Trending/ag/agno
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考