ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

agno 数据标注规模化(Scale-Out):用异步扇出、断点续跑与成本核算把单行标注循环升级到 10 万行

agno 数据标注规模化(Scale-Out):用异步扇出、断点续跑与成本核算把单行标注循环升级到 10 万行 agno 数据标注规模化Scale-Out用异步扇出、断点续跑与成本核算把单行标注循环升级到 10 万行【免费下载链接】agnoBuild, run, and manage agent platforms.项目地址: https://gitcode.com/GitHub_Trending/ag/agno本文围绕 agno cookbook 中data_labeling/_26_scale_out目录的完整实现展开讲解如何把本 cookbook 其他章节「同步循环逐行标注」的形态升级为面向真实数据集规模的异步 fan-out 管线。读完本文你将掌握三件事如何用asyncio.Semaphore控制并发并实测加速比而非口头宣称如何用 JSONL 检查点让 10 万行任务在被中断后只补跑未完成的行以及如何基于run.metrics在提交任务前就把 token 与美元成本算清楚。为什么需要单独的 Scale-Out 目录cookbook/data_labeling/下其他每个目录例如_01_text_classification/都在一个同步循环里标注十几行数据——这适用于演示但扛不住真实数据集按每行数秒的延迟估算10 万行意味着数天的墙钟时间。_26_scale_out的思路很克制标注调用本身完全不变——仍然是_01_text_classification/那套形态一个复用的 Agent、一个 Pydantic schema、每行一个标签。新增的全部是「管线骨架」带限流信号量的异步 fan-out让中断成本趋近于零的检查点在投入前就给任务定价的 token / 美元核算。也就是说这个目录是一套可移植的规模化外壳与具体标注任务解耦。文件清单与运行方式该目录共 4 个文件其中 3 个 Python 示例沿同一任务形态逐级叠加能力文件相对basic.py的新增能力basic.py纯异步 fan-out1 个复用 Agent 标注 30 条短评论asyncio.Semaphore(8)限流每 10 行打印进度resumable.py检查点续跑每完成一行立即追加写入data/generated/labels.jsonl重启时跳过已完成行with_cost_tracking.py成本核算从run.metrics聚合 token按列表价估算本运行成本并投影到 10 万行TEST_LOG.md2026-07-18 对gemini-3.5-flash、agno 2.7.4 的实测记录三个脚本全部 PASS运行方式三选一或依次运行python cookbook/data_labeling/_26_scale_out/basic.py python cookbook/data_labeling/_26_scale_out/resumable.py python cookbook/data_labeling/_26_scale_out/with_cost_tracking.py三个脚本均使用 Google Gemini 模型需要环境变量GOOGLE_API_KEY。basic.py把同步循环改造成异步扇出任务形态与文本分类完全一致basic.py复用了_01_text_classification的任务形态源码class Classification(BaseModel): label: Literal[positive, negative, neutral] Field( ..., descriptionThe assigned sentiment label ) TEXTS [ ... 30 条商品评论 ... ] ROWS [{id: fr{i:02d}, text: text} for i, text in enumerate(TEXTS, start1)] CONCURRENCY 8 PROGRESS_EVERY 10关键点每条数据被组织为带id的字典。这个id在resumable.py中会成为续跑的键因此从basic.py开始就保留了它。标注 Agent单实例复用temperature0labeler Agent( modelGemini(idgemini-3.5-flash, temperature0), instructionsYou classify product reviews by sentiment., output_schemaClassification, )单实例复用30 行共享同一个 Agent 对象通过await labeler.arun(row[text])逐行调用arun是 agno Agent 的异步入口源码位于 libs/agno/agno/agent/agent.py内部经由_run.arun_dispatch分发。temperature0让同一行在重跑时拿到相同标签。当然模型更新与服务端非确定性仍可能造成标签漂移这是任何 LLM 标注方案都无法完全消除的。带限流的并发执行与重试核心在label_row源码SEM asyncio.Semaphore(CONCURRENCY) async def label_row(row: dict, progress: Counter) - dict: async with SEM: start time.perf_counter() content None for attempt in range(3): # 重试 schema 断裂与瞬时 API 错误 try: run await labeler.arun(row[text]) except Exception: await asyncio.sleep(2**attempt) continue if isinstance(run.content, Classification): content run.content break if content is None: raise RuntimeError(frow {row[id]}: no valid label after 3 attempts) latency time.perf_counter() - start progress[done] 1 if progress[done] % PROGRESS_EVERY 0: print(flabeled {progress[done]}/{len(ROWS)} rows) return {id: row[id], text: row[text], label: content.label, latency: latency}值得注意的三个设计决策计时在信号量内部time.perf_counter()的起止都在async with SEM:之内。排队等待槽位的时间不计入行延迟——因为同步串行执行本来也不需要排队。这让后面的「顺序估算」与「墙钟时间」成为同一次运行的两个观测值。指数退避重试最多 3 次尝试失败后await asyncio.sleep(2**attempt)1s、2s同时校验返回内容是否真正是Classification实例——既兜住瞬时 API 错误也兜住 schema 断裂。进度用collections.Counter多个协程并发更新同一个 Counter在 CPython 下计数操作是原子的简单可靠。加速比是测出来的不是宣称的主流程用asyncio.gather并发调度全部 30 行源码results await asyncio.gather(*[label_row(row, progress) for row in ROWS]) wall_clock time.perf_counter() - wall_start ... mean_latency sum(latencies) / len(latencies) sequential_estimate mean_latency * len(results) print(fmeasured speedup: {sequential_estimate / wall_clock:.1f}x)输出示例来自 TEST_LOG.md 的实测注意数值随运行波动wall clock: 7.3s mean per-row latency: 1.72s sequential estimate: 30 rows x 1.72s 51.7s measured speedup: 7.0x在并发度 8 下实测加速约 7.0x标签分布与设计完全一致{positive: 10, negative: 10, neutral: 10}。由于是同一批请求竞争网络与服务端资源加速比不可能严格等于并发度这正说明「实测」的价值。resumable.pyJSONL 检查点让中断变廉价输出文件同时就是检查点resumable.py在basic.py之上只加了一样东西检查点源码。CHECKPOINT_PATH Path(__file__).parent / data / generated / labels.jsonl def load_done_ids(path: Path) - set: if not path.exists(): return set() with path.open() as f: return {json.loads(line)[id] for line in f if line.strip()}每一行标注完成的瞬间就追加写入并立即flush()checkpoint.write(json.dumps(result) \n) checkpoint.flush()flush()是关键asyncio.gather并发下只有把缓冲数据落到磁盘进程被 kill 时已完成的每一行才能存活。这样 10 万行任务跑到第 6 万行被杀掉重跑只做 4 万行的活。启动即跳过已完成行label_batch在启动时加载doneid 集合只调度未完成行done load_done_ids(CHECKPOINT_PATH) todo [row for row in rows if row[id] not in done] skipped len(rows) - len(todo)两段演示诚实地证明续跑main()先只给前 15 行模拟一次中断再给完整 30 行列表重跑await label_batch(ROWS[:15]) # pass 1: 模拟中断 await label_batch(ROWS) # pass 2: 全量续跑演示开始时删除旧检查点以保证确定性CHECKPOINT_PATH.unlink(missing_okTrue)。实测输出pass 1: wrote 15 rows, skipped 0 already labeled, checkpoint now has 15 pass 2: wrote 15 rows, skipped 15 already labeled, checkpoint now has 30重读labels.jsonl确认 30 行、30 个唯一 id、键为id/text/label、标签分布{positive:10, negative:10, neutral:10}。一个真实的坑事件循环绑定TEST_LOG 记录了一个值得注意的实现细节首个版本因两次asyncio.run()共享模块级信号量而报Semaphore is bound to a different event loop修复方式是把两段 pass 放进同一个asyncio.run(main())。如果你的任务也用模块级Semaphore务必让所有协程在同一个事件循环内创建与执行。with_cost_tracking.py提交前先定价数据来源RunOutput.metrics每一个RunOutput都携带run.metrics其中包含input_tokens、output_tokens、reasoning_tokens等字段。这些字段来自 agno 的RunMetrics数据类libs/agno/agno/metrics.py它继承自BaseMetrics定义了input_tokens、output_tokens、total_tokens、reasoning_tokens、cost等计数模型调用返回后由accumulate_model_metrics汇总进 run 级指标metrics.py。脚本在每行成功调用后提取这些值源码metrics run.metrics # 成功调用的 RunMetrics可能为 None return { ... input_tokens: metrics.input_tokens if metrics is not None else 0, output_tokens: metrics.output_tokens if metrics is not None else 0, reasoning_tokens: metrics.reasoning_tokens if metrics is not None else 0, has_metrics: metrics is not None, }定价常量写清楚价格口径价格以模块级常量明示保证结果可复算源码INPUT_PRICE_PER_1M 1.50 # Gemini 交互 tier 列表价美元 / 百万 token OUTPUT_PRICE_PER_1M 9.00 BATCH_DISCOUNT 0.5 # Batch API 为交互列表价的 50% PROJECTION_ROWS 100_000推理 token 才是账单大头gemini-3.5-flash是推理模型thinking token 在 metrics 中单独上报但按输出价计费所以可计费输出为billable_output total_output total_reasoning run_cost ( total_input / 1_000_000 * INPUT_PRICE_PER_1M billable_output / 1_000_000 * OUTPUT_PRICE_PER_1M )实测数据TEST_LOG.md指标总量每行均值input_tokens61320.4output_tokens1715.7reasoning_tokens4,459148.6推理 token 每行约 149 个而输出 token 仅约 6 个——思考 token 主导了账单。30/30 行都有 metrics。投影到 10 万行一个「耸肩」还是一条预算线estimated cost this run: $0.0426 ($1.420 per 1000 rows, interactive list prices) projected 100,000 rows: $141.97 interactive, $70.98 via the batch API这个投影就是标题里「定价任务」的意义在不敏感的批处理场景下Batch API 以交互列表价约 50% 的价格运行同一模型10 万行时就是 $142 与 $71 的差别。两点必要的诚实标注代码注释中已写明重试过的行实际计费高于其上报值因为只有成功的那次调用被计数token 数与成本随运行波动以上是单次观测。输出文件示例resumable.py写入的每行 JSON输出文件兼作检查点因此id即续跑键{id: r01, text: Absolutely love this blender, it crushes ice in seconds., label: positive} {id: r15, text: Returned it immediately, the fan noise is unbearable., label: negative} {id: r21, text: The box contains the charger, a cable, and a manual., label: neutral}何时使用这套外壳在真实数据集上运行任意标注任务这套管线从不窥探单行标注调用内部——把 schema 和 instructions 从其他目录换进来fan-out、检查点、成本核算原样生效_03_text_extraction/结构化字段抽取_15_document_classification/文档分类_17_llm_as_judge/LLM 裁判或任何同级目录。任务长到可能被中断崩溃、限流、合上笔记本——凡是有中断风险的场景用resumable.py。其续跑粒度精确到「行」不是「批次」。提交前先定价用with_cost_tracking.py算出每千行成本再决定是否开工当任务不敏感于延迟时优先评估 provider 的 Batch API约 50% 折扣。标注后的过滤与打包_22_dataset_curation/负责过滤、去重、打包已标注数据。它的 judge gate 与这里的每行调用形态一致因此同样可以套用本目录的外壳做规模化。与底层实现的对应关系RunOutput.metrics的类型即RunMetrics继承自BaseMetricsmetrics.py除本示例用到的三个 token 字段外还含total_tokens、cache_read_tokens、cache_write_tokens、cost等适合扩展更细粒度的成本统计。agent.arun的异步语义由 libs/agno/agno/agent/agent.py 的arun入口与arun_dispatch分发实现返回RunOutput非流式或异步迭代器流式。所有 token 计数最终由accumulate_model_metricsmetrics.py按(provider, id)聚合到details与顶层字段这与with_cost_tracking.py直接读run.metrics.input_tokens等字段的来源一致。小结_26_scale_out提供的是一个与任务解耦的规模化三层管线basic.py证明并发加速是可测量的并发 8 实测约 7.0xresumable.py用「边完成边落盘」的 JSONL 检查点把中断成本压到只重跑未完成行with_cost_tracking.py用run.metrics把账单算到每千行与 10 万行投影。把这三层外壳套在任意单行标注任务上就是一套从几十行演示走向真实数据集的完整路径。【免费下载链接】agnoBuild, run, and manage agent platforms.项目地址: https://gitcode.com/GitHub_Trending/ag/agno创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表