ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Agno 异步 Rollouts 实战:并发验证 Agent 并导出 SFT 数据集

2026/9/11 18:51:39 拓冰建站 浏览量
Agno 异步 Rollouts 实战:并发验证 Agent 并导出 SFT 数据集 Agno 异步 Rollouts 实战并发验证 Agent 并导出 SFT 数据集【免费下载链接】agnoBuild, run, and manage agent platforms.项目地址: https://gitcode.com/GitHub_Trending/ag/agno导读本文基于 cookbook/environments/_08_async_rollouts/ 目录下的示例与测试日志讲解如何在 Agno 的 Environments 体系中从异步应用中并发运行独立尝试rollouts先用arun_rollouts()以有界并发bounded concurrency批量执行任务并逐次评分再通过异步导出门面ato_sft_jsonl()将学习区learning zone内的通过轨迹导出为可直接喂给训练器的对话式 SFT JSONL。读完本文你将掌握异步 rollouts 的完整调用链、结果对象的使用方法、源码级的尝试隔离原理以及如何解读真实的测试通过率日志。一、异步 Rollouts同一个语义更短的墙钟时间Environments 是 Agno 用于Agent 验证与数据集生成的模块28 个渐进式示例目录共 79 个可运行单文件示例让 Agent 对一批困难任务各运行 K 次逐次评分检查通过率网格并导出通过的文字轨迹作为监督微调数据集参见 cookbook/environments/README.md。_08_async_rollouts这一节解决的问题是当你的应用本身已经跑在事件循环里异步服务、带异步入口的 notebook、批量任务时如何并发地跑这批独立尝试。与同步入口相比异步 rollouts 有一个明确而克制的承诺并发只改变墙钟时间wall-clock time不改变被记录尝试的含义或顺序。也就是说arun_rollouts()产出的结果形状与同步的run_rollouts()完全一致——包含网格grid、通过率、指纹fingerprints与学习区learning-zone选择。模型延迟占主导的场景下把 K×N 次独立调用并发出去能显著压缩总耗时但统计语义不变。典型适用场景见 README服务化场景Agent 验证逻辑内嵌于一个已经运行着 asyncio 事件循环的服务进程带异步入口的 Notebook不想为 rollouts 再开一个独立的同步线程批量任务一次要跑大量尝试模型延迟是瓶颈希望用并发把等待时间摊薄。在示例序列中_08_async_rollouts的位置也很明确先通过_07_difficulty_calibration/把任务难度校准到“强模型不再全绿”再到这里做并发的完整跑批后续的_09_task_selection/则只运行选定子集。与同步入口的关系从源码看同步入口只是异步入口的一层薄包装。在 runner.py 中def run_rollouts(env, *, k8, tasksNone, modelNone, concurrency4): try: asyncio.get_running_loop() except RuntimeError: pass else: raise RuntimeError(run_rollouts cannot be called from a running event loop; await arun_rollouts instead) return asyncio.run(arun_rollouts(env, kk, taskstasks, modelmodel, concurrencyconcurrency))注意这个守卫如果你已经在事件循环里直接调用run_rollouts()会抛出RuntimeError——此时应当await arun_rollouts()。这正是本目录存在的原因异步应用不能退出当前循环去跑同步版本。二、环境准备与运行所有 Environments 示例都使用OpenAIResponsesgpt-5.5模型需要设置OPENAI_API_KEY。从仓库根目录运行python cookbook/environments/_08_async_rollouts/basic.py python cookbook/environments/_08_async_rollouts/async_export.py两个文件都没有额外依赖直接可跑。需要明确的是README 与 Environments 总览均强调导出只生成数据集文件并不会训练模型本目录做的是“独立 rollout 并在完成后评分”不是实时的 RL 奖励循环。若想复用仓库自带的演示环境也可以先执行 scripts/demo_setup.sh再通过direnv exec . .venvs/demo/bin/python ...运行。三、basic.pyawaitarun_rollouts()跑并发 typed rolloutsbasic.py 是最小完整示例一个要求“精确计算”的 Agent两条边沿难度的大数算术任务用CodeScorer做类型化验证然后以k4, concurrency4并发跑完。3.1 类型化输出与评分器from pydantic import BaseModel, Field class FinalInteger(BaseModel): value: int Field(descriptionThe final integer after every requested operation) def exact_integer(run, expected) - bool: return isinstance(run.content, FinalInteger) and run.content.value expectedFinalInteger是 Agent 的输出 schema强制模型“只返回最终整数”且value的类型必须是intexact_integer是传给CodeScorer的评分函数它同时校验类型正确run.content是FinalInteger实例与数值精确等于期望值。这正是“typed rollouts”的含义——连类型错误都算失败。3.2 Agent 与任务定义agent Agent( modelOpenAIResponses(idgpt-5.5, reasoning_effortlow), instructionsCalculate exactly. Return only the final integer in the response schema., output_schemaFinalInteger, )任务使用了“精确算术阶梯”exact arithmetic ladders先做大整数乘法再要求对乘积做数字求和、乘系数、再取模。这类任务在单一运算上容易饱和叠加成阶梯后难度显著上升参见_21_math/的主题说明。env Environment( nameasync-rollouts, agentagent, tasks( Task(idasync-edge-a, inputMultiply 2718281828459045 by 1618033988749895. Add the decimal digits of the product, multiply that digit sum by 131071, then subtract the products remainder modulo 65521., expected20944939), Task(idasync-edge-b, inputMultiply 3141592653589793 by 1414213562373095. Add the decimal digits of the product, multiply that digit sum by 65537, then subtract the products remainder modulo 32749., expected10481347), ), scorerCodeScorer(exact_integer), )这里用到的三个核心类型定义于 environment.py类型字段/职责说明Taskinput、expected、id、metadata一行任务id缺省时在运行开始时按位置解析为t1..tN重复 id 会在Environment构造时被拒绝Environmentname、tasks、scorer、agent、timeout_seconds一次 rollouts 运行的单元timeout_seconds默认120agent可以是 Agent 实例或返回 Agent 的零参工厂CodeScorer接受一个评分函数对每一次尝试的run与expected打分值得注意Environment是frozenTrue的 dataclass字段不可重绑从结构上保证“结果确实来自这一组任务集、评分器与策略对象”。3.3 异步入口async def main() - None: result await arun_rollouts(env, k4, concurrency4) print(result) for task_result in result.task_results: print(f{task_result.task.id}: pass rate {task_result.pass_rate}) if __name__ __main__: asyncio.run(main())arun_rollouts的完整签名runner.pyasync def arun_rollouts( env: Environment, *, k: int 8, tasks: Optional[Sequence[Task]] None, model: Optional[Model] None, concurrency: int 4, ) - EnvironmentRunResult:参数默认值含义env必填要运行的Environment任务集 评分器 Agentk8每个任务运行多少次tasksNone可选子集选择必须取自env.tasks按身份选择重建的 task 不被接受None表示全部任务modelNone模型覆盖必须是Model实例而非字符串且会参与策略指纹计算concurrency4同时进行中的尝试数上限——即“有界并发”的界底层由 _engine.py 的arun_batch承担调度用asyncio.Semaphore限制并发窗口任务按input-major的顺序调度先铺开所有任务的第一批再填后续每个(input, attempt)对各自独立完成、独立评分。打印result会得到整个运行网格含每次尝试的评分与错误信息result.task_results中每个TaskResult暴露task.id、pass_rate、mean_value、in_learning_zone等属性见 runner.py。四、async_export.py异步验证 learning-zone 过滤 异步 SFT 导出async_export.py 在 basic.py 之上增加了一条完整的数据集生产流水线异步验证 → 学习区选择 → 只导出通过的文字轨迹。async def main() - None: result await arun_rollouts(env, k4, concurrency4) print(result) zone result.learning_zone() output_path Path(__file__).parent / data / generated / learning_zone.jsonl report await ato_sft_jsonl(zone, output_path) print(flearning-zone tasks: {len(zone.task_results)}) print(ftraining rows written: {report.n_written}) print(fdataset: {output_path})三步各司其职arun_rollouts(env, k4, concurrency4)与 basic.py 相同跑完整并发验证result.learning_zone()返回一个过滤副本只保留“既有通过的评分尝试、又有失败的评分尝试”的任务0 n_passed n_scored源码见 runner.py 与 runner.py。这些任务处于可学中间带策略有能力但不稳定每次失败都有一次通过作为对照——这正是 SFT 想要的样本区间await ato_sft_jsonl(zone, output_path)把学习区内任务的全部通过文字尝试写成对话式 SFT JSONL并返回ExportReport含n_written等计数。导出目录data/generated/由 exporter 自动创建。4.1 异步导出门面的实现ato_sft_jsonl是同步函数to_sft_jsonl的异步孪生sft.pyasync def ato_sft_jsonl(result, path, *, only_passedTrue): Async twin of to_sft_jsonl. return await asyncio.to_thread(to_sft_jsonl, result, path, only_passedonly_passed)实现非常朴素asyncio.to_thread把同步导出丢到线程池避免阻塞事件循环。这与 basic.py 中“同步评分器跑在线程里asyncio.to_thread”的处理保持一致——文件 IO 与 CPU 密集的 JSON 序列化都不应占用事件循环。4.2 导出格式与可移植性to_sft_jsonl输出的核心格式sft.py 的模块文档说明了设计动机每行一个 JSON 对象{messages: [{role: system|user|assistant, content: ...}]}Tinker、Together、Fireworks、OpenAI 都接受这个核心格式因此文件“通过省略而非翻译”获得可移植性——导出器刻意不写tools、weight、trainable键因为最严格的消费方做集合相等校验任何额外键都会导致整个文件被拒绝导出内容的两个关键约束系统消息被保留它承载了诱导输出的格式指令assistant 文本取run.messages中最后一条 assistant 消息的原始content逐字节原文绝不重新序列化run.content——因为output_schema下的run.content是 Pydantic 模型它的str()/model_dump_json()产出的文本并非模型真正生成的字符串见 sft.py。4.3 跳过规则与文件上限每个尝试按顺序经过如下判定_classify见 sft.py每条规则对应ExportReport中的一个计数规则计数理由未评分score is None不进入计数非候选永不进入导出器only_passedTrue且未通过n_skipped_failed失败答案不能写进监督文件触发了工具调用上限n_skipped_limit_hit被拒绝工具调用后给出的答案是“胁迫下的答案”运行中使用了工具n_skipped_tool_runs交集格式没有工具表示只导最终答案会让模型学会“不用工具作答”无导出的文字对话n_skipped_no_text没有可用的 assistant 文本文件级上限每条文件最多 320 条对话、1 MiB超限行按发射顺序从尾部丢弃并计入n_dropped_over_cap绝不静默截断。写入时显式使用newline关闭平台换行转换避免 Windows 的 CRLF 把恰好压线文件推过字节上限、破坏确定性见 sft.py。4.4 可审计的 sidecar除 JSONL 本体外导出器还会在同路径写一份learning_zone.jsonl.meta.jsonsidecar携带env_fingerprint、policy_fingerprint、六项report计数、only_passed选项以及每行数据的来源task_id、attempt_index、score。指纹让每一行训练数据都能追溯到产生它的环境与策略这正是“可审计训练集”的落点相关内容可继续阅读_11_export_provenance/。五、源码级原理尝试隔离、指纹与错误风暴5.1 无条件尝试隔离arun_rollouts最核心的设计是隔离不可关闭isolation is unconditional and there is no knob。每一次尝试都运行在全新的内存数据库fresh in-memory store全新的 session 与 user id关闭响应缓存response cache off的 Agent 副本之上。随后生产环境的解析器原封不动地在这份输入上运行——因此尝试的 prompt 与一个全新生产用户看到的完全一致见 runner.py。隔离只切断写路径关闭update_memory_on_run运行后的记忆抽取调用关闭enable_agentic_memoryupdate_user_memory 工具及其 prompt 块关闭update_knowledgeknowledge 写工具关闭enable_session_summaries运行后的摘要写一次额外的 LLM 调用置空save_response_to_file否则 K 次尝试会竞争写入调用者的文件。读路径被保留knowledge 检索走knowledge.vector_db而非agent.db因此带 RAG 的 Agent 在 rollout 中能正常检索记忆读取则解析到尝试自己的空库渲染出“全新用户”应有的空状态——调用者世界的用户态数据绝不能泄漏进样本。仓库测试 test_runner.py 用Recorder快照验证了“每次尝试都有独立内存库”“knowledge/memory/learning 写全被关闭”等约束。5.2 双重指纹环境指纹与策略指纹每次运行开始时计算两个 sha256 指纹并打在结果上runner.pyenv_fingerprint前缀envfp2:对任务列表、评分器摘要、声明的工具 schema、tool_choice、所有 prompt 成形字段与上下文标志、模型 prompt 载荷、session_state、终止设置做哈希environment.pypolicy_fingerprint对模型类、id、provider、base_url 及所有请求塑形参数做哈希environment.py。两者分离的价值模型 id如gpt-5.5vsgpt-5.5-mini不会哈希成同一值而“prompt 改动”算环境变化、“模型采样参数改动”算策略变化。任何组件无法指纹化时降级为None并告警绝不崩掉整个运行。5.3 错误风暴中止error-storm运行会监控开头的一批尝试若前max(concurrency, 4)次尝试全部以同类错误结束同error_type判定为统一配置错误坏 key、不可达的 base_url随即停止调度后续尝试并在结果上标注stopped_earlyerror-stormrunner.py。该检查只对单任务运行生效——多任务下调度是 input-major 的开头的一批可能都来自同一个倒霉任务此时中止会误伤健康任务。5.4 三个诚实的残留源码明确列出 override 触及不到的三类情况runner.py同步评分器不可中断env.timeout_seconds只约束尝试“返回”的时间同步评分器跑在线程中asyncio.to_thread取消无法打断它被放弃的线程会自然跑完Tracing 是进程全局状态配置了setup_tracing时尝试的 trace/span 行会像普通运行一样写入调用者的 trace 存储可通过rollout-*前缀的 session/user id 识别用户提供的 callables 共享引用pre/post hooks 每次尝试都会触发fallback_config.callback也会在尝试回退时触发——让 rollout 可见的 hooks 保持幂等或按它们收到的rollout-*id 过滤。六、使用约束与注意点MCP 工具必须以工厂方式传入如果env.agent直接持有MCPTools或藏在 callable tools 工厂/工具集后面运行会在任何尝试开始前被拒绝——并发尝试共享一个 MCP session 会在中途互相连接/关闭导致整批丢失工厂里每次调用新建MCP 实例才是正确姿势runner.py。模型必须是 Model 实例model参数不接受字符串字符串解析被刻意禁止。子集选择按身份tasks必须从env.tasks中按身份选取重建的 task 会被拒绝。工具型任务可验证但不可导出当前 SFT 格式只支持文字轨迹带工具的 run 会被导出器跳过n_skipped_tool_runs而不是丢掉工具证据去教模型“无工具作答”。only_passed是拒绝采样learning_zone()选任务、only_passed选任务内的尝试二者配合才得到正确的 SFT 输入。直接导出完整结果会按 K 倍加权饱和任务属于已知的选择效应sft.py。七、测试日志解读一份真实的通过率记录TEST_LOG.md 记录了本目录在2026-07-20、gpt-5.5经OpenAIResponses、Agno 2.7.4下的真实运行结果是理解“异步 rollouts 到底输出什么”的第一手证据。basic.py —— PASSConcurrent typed rollouts througharun_rollouts().async-edge-a和async-edge-b各通过 3/40.75。所有八次尝试都被评分且两行都在 true binary learning zone。解读k4意味着每条任务 4 次尝试两条任务共 8 次全部被评分无超时、无 scorer 崩溃两条任务的pass_rate均为 0.75满足0 n_passed n_scored因此都落在二元学习区——正是learning_zone()会保留、SFT 导出会消费的“有能力但不稳定”的中间带这与“任务难度已校准”的设计意图吻合校准的目标就是让强模型不要整墙全绿参见 cookbook/environments/README.md。async_export.py —— PASSAsync verification followed by passing-only SFT JSONL export. 最终异步导出运行中export-edge-a通过 1/40.25export-edge-b通过 4/41.00。学习区选择保留了第一个任务ato_sft_jsonl()写入了它的那一条通过对话。解读两条任务的通过率差异很大0.25 vs 1.00这是单次小样本运行的自然波动也解释了为什么 rollouts 要做多次尝试取分布export-edge-b4/4 全绿、export-edge-a1/4 部分通过 →learning_zone()只保留export-edge-a它既有失败也有通过该任务的 4 次尝试中只有 1 次通过且无工具调用、有可用文本 →ato_sft_jsonl恰好写出1 行训练数据n_written 1n_skipped_failed 3。日志还记录了一次有价值的对照在把导出调用切换为异步孪生前的 pre-final 实测运行中观察到 4/4 与 3/4并写入了三行上文最终通过率对应的是检入仓库的异步代码路径。这段记录说明两件事其一同一个环境、同样的k不同批次的通过率会波动——评估结论必须建立在多次运行与分布之上其二异步路径arun_rolloutsato_sft_jsonl与同步路径run_rolloutsto_sft_jsonl语义等价最终日志以异步代码路径的结果为准。八、从验证到数据集完整的落地链路把_08_async_rollouts放进 Environments 的整体工作流中一条可复用的链路是校准用_07_difficulty_calibration/调整任务难度直到强模型不再全绿异步并发验证本文的basic.py——await arun_rollouts(env, k..., concurrency...)拿到网格与通过率学习区选择result.learning_zone()筛出可学中间带异步导出本文的async_export.py——await ato_sft_jsonl(zone, path)只导出通过的文字轨迹同时落一份带指纹的.meta.json溯源 sidecar子集复跑需要补证据时用_09_task_selection/只跑选定子集。而本目录的 TEST_LOG 则是这条链路的“验收记录”它证明异步路径在真实模型上可复现地输出与同步路径一致的结果形状并给出了可对照的通过率基线0.75 / 0.251.00。对于任何要在异步服务或批量流水线里做 Agent 验证与数据生产的场景arun_rolloutsato_sft_jsonl就是 Agno 提供的开箱即用方案。相关源码与文档索引示例源码cookbook/environments/_08_async_rollouts/basic.py、cookbook/environments/_08_async_rollouts/async_export.py运行器与隔离逻辑libs/agno/agno/environments/runner.py、并发批处理 libs/agno/agno/environments/_engine.py环境与任务类型、指纹实现libs/agno/agno/environments/environment.pySFT 导出器libs/agno/agno/environments/exporters/sft.py隔离约束测试libs/agno/tests/unit/environments/test_runner.pyEnvironments 总览cookbook/environments/README.md【免费下载链接】agnoBuild, run, and manage agent platforms.项目地址: https://gitcode.com/GitHub_Trending/ag/agno创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考