资讯动态

Mastra 并行工作流实战:用 .parallel() 构建高性能并行分析与结果汇总工作流

发布时间:2026/9/14 18:57:18 来源:尧图企业网站定制
Mastra 并行工作流实战用 .parallel() 构建高性能并行分析与结果汇总工作流【免费下载链接】mastraMastra is the modern TypeScript framework for AI-powered applications and agents.项目地址: https://gitcode.com/GitHub_Trending/ma/mastra本文以 Mastra 工作流教程「Building Parallel Workflow」一课为核心讲解如何使用createWorkflow().parallel()让多个互不依赖的分析步骤SEO 分析、可读性分析、情感分析同时执行并通过一个.then()汇总步骤将各分支结果聚合为结构化输出。读完后你将掌握 Mastra 中并行执行的数据流机制、结果对象的键名规则以及parallel()方法在源码层面的实现原理能够为耗时型多分析场景网络请求、AI 调用、重计算设计出真正提速的工作流。前置背景三个可独立执行的并行步骤在构建并行工作流之前教程的上一节Creating Parallel Steps已经定义了三个createStep步骤它们共享同一份输入content与type但彼此之间没有任何数据依赖SEO 分析步骤id: seo-analysis输出seoScore60~99 的随机整数和keywords内容中长度大于 4 的单词前 3 个模拟耗时 800ms可读性分析步骤id: readability-analysis基于句均词数计算readabilityScore100 - avgWordsPerSentence * 3下限 0与gradeLevelEasy/Medium/Hard三档模拟耗时 600ms情感分析步骤id: sentiment-analysis通过统计正负关键词命中次数得出sentimentpositive/neutral/negative与confidence0.7~1.0模拟耗时 700ms。这三个步骤正是典型的并行适用场景符合 Understanding Parallel Execution 中给出的三个判断标准步骤互不依赖、执行耗时模拟中为 600~800ms真实场景多为网络请求与 AI 调用、处理同一份输入。如果串行执行总耗时约为 800 600 700 2100ms而并行执行的总耗时只取决于最慢的分支约为 800ms。创建并行工作流在已有的三个步骤基础上用createWorkflow定义并行工作流。以下代码完整继承了教程中的写法可直接加入你的工作流文件export const parallelAnalysisWorkflow createWorkflow({ id: parallel-analysis-workflow, description: Run multiple content analyses in parallel, inputSchema: z.object({ content: z.string(), type: z.enum([article, blog, social]).default(article), }), outputSchema: z.object({ results: z.object({ seo: z.object({ seoScore: z.number(), keywords: z.array(z.string()), }), readability: z.object({ readabilityScore: z.number(), gradeLevel: z.string(), }), sentiment: z.object({ sentiment: z.enum([positive, neutral, negative]), confidence: z.number(), }), }), }), }) .parallel([seoAnalysisStep, readabilityStep, sentimentStep]) .then( createStep({ id: combine-results, description: Combines parallel analysis results, inputSchema: z.object({ seo-analysis: z.object({ seoScore: z.number(), keywords: z.array(z.string()), }), readability-analysis: z.object({ readabilityScore: z.number(), gradeLevel: z.string(), }), sentiment-analysis: z.object({ sentiment: z.enum([positive, neutral, negative]), confidence: z.number(), }), }), outputSchema: z.object({ results: z.object({ seo: z.object({ seoScore: z.number(), keywords: z.array(z.string()), }), readability: z.object({ readabilityScore: z.number(), gradeLevel: z.string(), }), sentiment: z.object({ sentiment: z.enum([positive, neutral, negative]), confidence: z.number(), }), }), }), execute: async ({ inputData }) { console.log( Combining parallel results...) return { results: { seo: inputData[seo-analysis], readability: inputData[readability-analysis], sentiment: inputData[sentiment-analysis], }, } }, }), ) .commit()这段代码的结构值得逐层拆解工作流入输出 SchemainputSchema与三个并行步骤的inputSchema完全一致content必填type缺省为articleoutputSchema则描述最终汇总后的结构即一个results对象内含seo、readability、sentiment三个字段。.parallel([step1, step2, step3])把三个步骤作为数组传入它们将同时接收工作流的输入数据并并发执行。.then(combineStep)并行块之后的下一步其inputSchema的键是各步骤的 step IDseo-analysis、readability-analysis、sentiment-analysis而不是步骤在数组中的位置——这是并行数据流的核心约定。.commit()冻结工作流定义生成可被注册到Mastra实例并运行/序列化的成品。并行数据流的四个阶段教程对.parallel()的数据流给出了四步描述这四点也是排查并行工作流问题的基本思路每个并行步骤收到相同的输入数据——即工作流入参本例为{ content, type }三个步骤的inputSchema都与之对齐这是parallel()能通过类型检查的前提步骤同时执行——三个步骤的execute函数并发运行互不阻塞结果被收集进一个以 step ID 为键的对象——并行块完成后的中间态形如{ seo-analysis: {...}, readability-analysis: {...}, sentiment-analysis: {...} }下一个步骤收到全部并行结果——汇总步骤的inputData就是这个对象execute中通过inputData[seo-analysis]这样的方式按 ID 取值重命名为seo/readability/sentiment后返回。源码级验证.parallel() 在 Mastra 中如何工作以上数据流约定并非只是文档描述可以直接在核心包源码中得到印证。Workflow类的parallel()方法定义在 workflow.ts其关键行为有三点记录并行节点方法调用this.stepFlow.push({ type: parallel, steps: [...] })把数组中的每个步骤统一转换为 step entry 后存入该节点第 2380~2389 行。也就是说工作流的执行拓扑里并行块是一个type: parallel的容器节点内部按数组顺序登记各分支。执行阶段对case parallel的分支处理也位于同一文件workflow.ts 第 1885 行附近。按 step ID 注册步骤steps.forEach(step { this.steps[step.id] step })第 2390~2392 行说明每个并行分支都在工作流内部以 ID 为键登记这正是「结果对象以 step ID 为键」这一约定的底层来源——step ID 是并行分支结果取值的唯一寻址方式。输出类型推导parallel()返回的新工作流类型中MappedOutputSchema被推导为{ [K in keyof StepsRecordTParallelSteps]: InferStandardSchemaOutput...[outputSchema] }第 2400~2404 行即「以各步骤 ID 为键、以各步骤outputSchema推导结果为值」的 Record。这意味着汇总步骤的inputSchema写成 step ID 形式不仅是运行时约定也是 TypeScript 层面的类型约束——如果combineStep的inputSchema键名与步骤 ID 不一致类型推导会直接报错。另外从 builder/index.ts 的源码结构看parallel与conditional、foreach、loop同属容器型 entryOPTIONAL_ID_ENTRY_TYPES容器节点本身的id是可选的而其内部每个单步骤必须声明 ID——这进一步解释了为什么本例中三个步骤的idseo-analysis等是必填且必须与汇总步骤输入键严格一致。性能收益与适用前提教程给出的量化预期是三个步骤模拟耗时分别为 800msSEO、600ms可读性、700ms情感串行总计约 2.1~2.2 秒并行后总耗时约等于最慢分支的 800ms。需要说明两点适用前提这个收益成立的前提是分支之间确实没有数据依赖且能并发消费资源网络并发、异步 IO。若步骤内部是纯同步 CPU 重计算受单线程限制收益会打折本例各分支的输入完全相同同一份content。Mastra 的parallel()将上游状态整体传给每个分支各分支可在自己的inputSchema中声明所需子集适合「同一份数据多种分析」的模式。教程后续会在 Testing Parallel Performance 一节中实际运行该工作流验证日志与耗时确实呈现并行特征三条console.log交替出现、总耗时约 800ms。你也可以参考核心包的并行相关测试用例例如 parallel-writer.test.ts 与 parallel-nested-restart.test.ts了解并行块在序列化、嵌套与恢复场景下的行为。关键要点小结要点说明.parallel([step1, step2, step3])让数组内所有步骤同时执行接收相同的上游输入结果对象的键使用各步骤的step ID如seo-analysis这是按源码中this.steps[step.id]注册机制推导出的取值规则汇总步骤在.parallel()之后用.then()接一个 combine 步骤一次性处理所有并行结果并整形为工作流最终输出类型保障parallel()的返回类型把后续输入推导为「step ID → 各步骤 outputSchema」的 Record键名写错会在编译期暴露按上述方式构建并行工作流后即可按教程进入测试环节观察并行执行下的性能提升。【免费下载链接】mastraMastra is the modern TypeScript framework for AI-powered applications and agents.项目地址: https://gitcode.com/GitHub_Trending/ma/mastra创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价