Agno Workflow 后台执行实战:异步轮询与 WebSocket 实时事件流 📅 发布时间:2026/9/9 19:55:15 👁 浏览次数: 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 放到后台运行backgroundTrue并通过两种方式取回结果一是非流式模式下基于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在常规入参之外额外暴露了四个与后台执行直接相关的参数参数类型作用backgroundOptional[bool] False是否以后台方式启动运行streamOptional[bool] None是否流式返回内容stream_eventsOptional[bool] None是否同时推送运行事件RunStarted/StepStarted 等websocketOptional[WebSocket]显式传入 WebSocket 连接同时会开启enable_websocketenable_websocketbool False与backgroundstream同时为 True 时改用 WebSocket 传输代替默认 SSE 传输源码中backgroundTrue时会按以下优先级路由到三种底层实现backgroundTruestreamTrue WebSocket 启用传了websocket或enable_websocketTrue→_arun_background_stream_ws实时事件 WebSocket 传输。代码中注释明确写着 Background Streaming WebSocket Real-time events (opt-in)backgroundTruestreamTrue未启用 WebSocket→_arun_background_stream后台 默认 SSE 传输backgroundTrue 非流式 →_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.pybackground_poll.py 的场景非常典型研究 Hacker News 与 Web 上的科技话题再由内容策划 Agent 输出为期四周的内容排期。整个过程耗时较长因此示例选择「异步发起 定时轮询」而不是同步阻塞等待。3.1 组装 Agent、Team 与 Step示例先创建三个 AgentHackernews Agent工具HackerNewsTools、Web Agent工具WebSearchTools负责研究Content Planner按指令规划内容排期再把两个研究 Agent 放进一个Team(nameResearch Team)。随后把两个执行单元各包成一个 Stepresearch_step Step(nameResearch Step, teamresearch_team) content_planning_step Step(nameContent Planning Step, agentcontent_planner)Step与Workflow的类型定义位于agno.workflow.step与agno.workflow.workflow。这里的核心思路是Agent 是单一能力执行者Team 是横向协作单元Step 是 Workflow 的最小编排节点——Step 既可挂 Agent 也可挂 Team编排层只需关心 Step 的先后顺序。3.2 用 SqliteDb 持久化会话Workflow 构造时传入了一个SqliteDbSQLite 会话数据库演示环境的库也覆盖在 cookbook/06_storage/sqlite 等示例中content_creation_workflow Workflow( nameContent Creation Workflow, descriptionAutomated content creation from blog posts to social media, dbSqliteDb( session_tableworkflow_session, db_filetmp/workflow.db, ), steps[research_step, content_planning_step], )配置项含义session_table会话表名这里用workflow_sessiondb_fileSQLite 文件落盘路径tmp/workflow.db会在运行时创建。为什么后台模式必须配 db后台任务发起后当前协程无法直接拿到最终结果只能通过run_id反查而反查的数据源正是会话数据库。没有 db就无法跨请求/跨进程恢复运行状态。3.3 发起后台任务并立即拿到 Initial Responsebg_response await content_creation_workflow.arun( inputAI trends in 2024, backgroundTrue, ) print(fInitial Response: {bg_response.status} - {bg_response.content}) print(fRun ID: {bg_response.run_id})backgroundTrue时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(fTimeout after {poll_count} attempts) break await asyncio.sleep(5) continue if result.has_completed(): # 运行完成的判据 break if poll_count 200: # 总轮询上限约 1000 秒 print(fTimeout 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, markdownTrue)美化打印——它会把 run_id、会话信息与各步最终内容以 Markdown 形式渲染出来。整个文件以asyncio.run(main())驱动说明这套轮询逻辑天然适配 FastAPI/AgentOS 等异步运行环境。四、方案二服务端——后台工作流事件经 WebSocket 推送websocket_server.pywebsocket_server.py 基于FastAPI uvicorn把「后台运行 Workflow 事件流式推送」封装成一个可被任意客户端连接的服务。4.1 服务拓扑与启动启动后监听0.0.0.0:8000WebSocket 端点ws://localhost:8000/wsHTTP 状态端点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( inputmessage, session_idsession_id, streamTrue, stream_eventsTrue, backgroundTrue, websocketwebsocket, )对照第二节的路由逻辑这一调用组合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.pywebsocket_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发送心跳并观察pongquit/exit/q退出并清理监听任务、断开连接未认证时连接提示栏会高亮显示 AUTHENTICATION REQUIRED提醒先输入auth。5.3 事件协议解析JSON 与 SSE 双格式兼容listen_for_events的解析策略体现了对两种服务端实现风格的兼容先尝试json.loads整体解析纯 JSON 事件如方案二服务端发送的消息解析失败则按SSE 文本格式二次解析形如event: Xdata: {...}的多行消息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.pynormal 模式执行 35 秒后超时日志显示已完成多轮 Agent Run说明完整跑完需要真实模型调用与更长执行窗口评估超时阈值时应放宽websocket_client.pystartup 模式通过8 秒内完成启动校验因当时无服务端而报连接失败Connect call failed属预期行为——它必须在服务端存活时才能完整演示websocket_server.pystartup 模式通过8 秒后按预期终止进程说明服务可正常拉起。七、把方案落地到生产的关键清单综合源码实现与示例细节在真实项目中落地「后台执行」时可沉淀如下经验选对后台模式只需要最终结果 → 非流式backgroundTrueget_run()轮询需要实时进度 →backgroundTrue streamTrue stream_eventsTrue并决定用默认 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),仅供参考