PocketFlow 并行批量翻译实战:用 AsyncFlow 与 AsyncParallelBatchNode 将多语言翻译耗时从 1136 秒压到 209 秒
PocketFlow 并行批量翻译实战用 AsyncFlow 与 AsyncParallelBatchNode 将多语言翻译耗时从 1136 秒压到 209 秒【免费下载链接】PocketFlowPocket Flow: 100-line LLM framework. Let Agents build Agents!项目地址: https://gitcode.com/gh_mirrors/poc/PocketFlow本文围绕 PocketFlow 仓库中的 parallel-batch 示例cookbook/pocketflow-parallel-batch讲解如何用AsyncFlow与AsyncParallelBatchNode将一份 Markdown 文档一次性并发翻译成 8 种语言并与串行BatchNode方案对比验证提速效果。读完本文你将掌握 PocketFlow 异步并行节点的三阶段编程模型prep_async/exec_async/post_async、asyncio.gather的底层并发原理以及一套可直接复用的批量 I/O 密集型任务并行化实战模板。一、场景与目标多语言翻译为什么必须并行将一份文档翻译成多种语言是典型的I/O 密集型批量任务每个目标语言都需要一次独立的 LLM API 调用而每次调用的大部分时间都花在网络往返与模型生成上CPU 几乎不参与。如果采用串行方式见 cookbook/pocketflow-batch 中的TranslateTextNodeFlow实现8 次翻译必须依次等待总耗时等于各次翻译耗时之和而并行方式下8 次请求同时发出总耗时约等于最慢的那一次。本示例的具体目标来自 关联文档读取仓库根目录的 README.md 作为源文本将其并行翻译为 8 种语言中文、西班牙语、日语、德语、俄语、葡萄牙语、法语、韩语每种语言分别保存为translations/README_LANGUAGE.md记录总耗时与串行版本对比量化并行带来的收益。二、环境准备与运行步骤示例位于 cookbook/pocketflow-parallel-batch/依赖pocketflow、anthropic、python-dotenv、httpx与aiofiles见 requirements.txt。第 1 步安装依赖pip install -r requirements.txt第 2 步设置 Anthropic API Keyexport ANTHROPIC_API_KEYyour-api-key-here将your-api-key-here替换为真实密钥也可以把ANTHROPIC_API_KEYyour-api-key-here写入.env文件依赖列表中的python-dotenv用于加载该文件。第 3 步可选验证 API Keypython utils.py该脚本会向模型发送一个简短的测试提示In a few words, what is the meaning of life?并打印响应用于确认密钥有效需先设置好密钥。第 4 步运行并行翻译python main.py程序会自动读取../../README.md创建translations/输出目录并发执行 8 个语言的翻译并保存结果。三、核心实现解析TranslateTextNodeParallel 节点整个并行逻辑收敛在一个继承自AsyncParallelBatchNode的节点上入口见 main.py。3.1 prep_async构造批量数据async def prep_async(self, shared): text shared.get(text, (No text provided)) languages shared.get(languages, []) return [(text, lang) for lang in languages]prep_async从共享存储shared中取出源文本与目标语言列表返回一个批次列表每个元素是(text, language)元组对应一次独立的翻译任务。这就是批量Batch语义的来源——批次的粒度决定了并发任务的数量。3.2 exec_async并发执行 LLM 调用async def exec_async(self, data_tuple): text, language data_tuple prompt f Please translate the following markdown file into {language}. But keep the original markdown format, links and code blocks. Directly return the translated text, without any other text or comments. Original: {text} Translated: result await call_llm(prompt) print(fTranslated {language} text) return {language: language, translation: result}exec_async是每个批次项的执行体构造翻译提示词await异步 LLM 调用返回{language: ..., translation: ...}字典。注意提示词中明确要求保留原始 Markdown 格式、链接与代码块只返回译文这是保证译文可直接落盘为.md文件的关键。call_llm来自 utils.py是基于AsyncAnthropic的异步封装async def call_llm(prompt): client AsyncAnthropic(api_keyos.environ.get(ANTHROPIC_API_KEY, your-api-key)) response await client.messages.create( modelclaude-3-7-sonnet-20250219, max_tokens20000, thinking{type: enabled, budget_tokens: 16000}, messages[{role: user, content: prompt}], ) return response.content[1].text值得注意的实现细节该封装启用了 Claude 的 thinking 模式budget_tokens16000因此响应内容列表的第一项是思考过程、第二项才是正文代码用response.content[1].text取正文。max_tokens20000的设定也说明它面向的是整篇 README 这样的大段文本翻译。3.3 post_async汇总结果并异步写盘async def post_async(self, shared, prep_res, exec_res_list): output_dir shared.get(output_dir, translations) os.makedirs(output_dir, exist_okTrue) for result in exec_res_list: if isinstance(result, dict): language result.get(language, unknown) translation result.get(translation, ) filename os.path.join(output_dir, fREADME_{language.upper()}.md) try: import aiofiles async with aiofiles.open(filename, w, encodingutf-8) as f: await f.write(translation) print(fSaved translation to {filename}) except ImportError: with open(filename, w, encodingutf-8) as f: f.write(translation) print(fSaved translation to {filename} (sync fallback)) except Exception as e: print(fError writing file {filename}: {e}) else: print(fWarning: Skipping invalid result item: {result}) return defaultpost_async接收全部并发任务的结果列表exec_res_list按language.upper()生成文件名如README_CHINESE.md优先使用aiofiles异步写盘若环境未安装aiofiles则回退到同步写入并对非字典结果做容错跳过。最后返回default作为动作信号单节点流程下即结束。3.4 组装流程并运行def create_parallel_translation_flow(): translate_node TranslateTextNodeParallel(max_retries3) return AsyncFlow(starttranslate_node)节点以max_retries3初始化随后被包进AsyncFlow。主函数中通过asyncio.run(main())启动核心是shared { text: text, languages: [Chinese, Spanish, Japanese, German, Russian, Portuguese, French, Korean], output_dir: translations } start_time time.perf_counter() await translation_flow.run_async(shared) end_time time.perf_counter() print(f\nTotal parallel translation time: {duration:.4f} seconds)用time.perf_counter()精确计时包裹run_async8 种语言与输出目录均通过shared共享存储注入体现了 PocketFlow 节点间通过共享字典传参的设计。四、源码级原理AsyncParallelBatchNode 如何实现并发AsyncParallelBatchNode的核心实现位于 pocketflow/init.py只有一行class AsyncParallelBatchNode(AsyncNode, BatchNode): async def _exec(self, items): return await asyncio.gather(*(super(AsyncParallelBatchNode, self)._exec(i) for i in items))它的工作方式是继承BatchNode因此拥有批次列表逐项执行的批量语义用asyncio.gather将每个批次项即每个(text, language)元组的exec_async任务同时调度而不是像AsyncBatchNodepocketflow/init.py那样[await ... for i in items]逐个等待。这就是串行批量与并行批量在框架层面的本质区别前者在循环中逐个await后者一次性发起所有协程并等待全部完成。与之配套的还有异步重试机制AsyncNode._execpocketflow/init.py会按max_retries循环捕获异常后若未达上限则await asyncio.sleep(wait)等待重试这也是示例中max_retries3能提高 LLM 调用稳定性的原因。异步流程编排则由AsyncFlow._orch_asyncpocketflow/init.py负责复制起始节点副本、循环执行prep_async → exec_async → post_async、依据返回值查找后继节点直到流程结束。本示例是单节点流程运行一次节点即完成全部 8 个翻译任务。五、串行 vs 并行示例运行对比关联文档给出了同一份 README、同样 8 种语言的两次实测输出对比# --- 串行运行输出来自 pocketflow-batch--- Starting sequential translation into 8 languages... Translated Chinese text ... Translated Korean text Saved translation to translations/README_CHINESE.md ... Saved translation to translations/README_KOREAN.md Total sequential translation time: ~1136 seconds # --- 并行运行输出本示例--- Starting parallel translation into 8 languages... Translated French text Translated Portuguese text ... # 各语言消息可能交错出现 Translated Spanish text Saved translation to translations/README_CHINESE.md ... Saved translation to translations/README_KOREAN.md Total parallel translation time: ~209 seconds两组关键观察耗时差异串行约 1136 秒并行约 209 秒本次运行约 5.4 倍加速。文档同时注明实际时间取决于 API 响应速度与系统环境加速倍数会随网络状况、模型负载、语言数量变化但数量级上的收益是稳定的——并行总耗时趋近于最慢单次翻译的耗时而非所有翻译之和。输出特征串行版本中 Translated X text 严格按语言顺序出现并行版本中各语言的打印消息交错出现这正是多个协程并发推进、各自在不同时刻完成的直观证据。最终落盘的文件则一致地汇聚到translations/目录本仓库中可见 translations/ 下的 8 个语言文件。六、为什么异步并行对 LLM 批量任务特别有效PocketFlow 的 parallel-batch-flow 示例 在图像处理场景下同样给出了量化的对比约 8 倍加速并总结了通用规律串行总时间 所有子任务时间之和。适合需要保持顺序、或受限于 API 速率限制的场景并行总时间 ≈ 最慢单个子任务的时间。适合 I/O 密集型、相互独立的任务。LLM 翻译恰好同时满足后者的两个前提每次 API 调用是独立的翻译中文与翻译韩语互不依赖且耗时几乎全部花在等待网络响应上I/O 密集。因此用AsyncParallelBatchNode替换BatchNode仅靠asyncio.gather并发调度即可获得近乎线性的收益无需引入多进程或多线程带来的额外内存与上下文切换开销。七、文件结构与扩展建议本示例的文件组织源自 关联文档 的 Files 一节文件作用main.py并行批量翻译节点与流程的实现含入口main()utils.py基于AsyncAnthropic的异步 LLM 调用封装requirements.txt项目依赖含aiofilestranslations/输出目录运行时自动创建在此基础上可做的实战扩展调整语言集合修改main.py中shared[languages]列表即可增减目标语言并发度随之自动变化复用同一模式凡是一个输入 × 多组独立参数的 I/O 密集型任务多语言摘要、多风格改写、批量文档分类、多仓库代码评审都可以照搬prep_async → exec_async → post_async三阶段结构控制并发上限若担心同时发起过多请求触发限流可在exec_async中引入asyncio.Semaphore限制并发窗口对比验证如需复现串行基线直接对照 cookbook/pocketflow-batch/main.py 中基于BatchNodeFlow的实现两者除节点基类与run_async/run调用外提示词与落盘逻辑几乎一致非常适合做 A/B 计时实验。总而言之这个示例是理解 PocketFlow 异步并行抽象的最佳入口它用最小的代码量同时展示了AsyncFlow、AsyncParallelBatchNode、asyncio.gather的协作方式以及并行化 I/O 密集任务总耗时从求和变为取最大值这条可迁移到任何 LLM 批量场景的核心经验。【免费下载链接】PocketFlowPocket Flow: 100-line LLM framework. Let Agents build Agents!项目地址: https://gitcode.com/gh_mirrors/poc/PocketFlow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考