第 8 章 · 分支与并行执行
本章目标:掌握 .branch() 条件分支与 .parallel() 并行扇出,学会在汇合点聚合多路结果。
8.1 超越线性:三种控制流
真实业务流程很少是一条直线。Mastra Workflow 提供三种控制流原语:
text
.then() 线性串行: A → B → C
.branch() 条件分支: A → [条件1→B / 条件2→C]
.parallel() 并行扇出: A → (B ∥ C) → 汇合8.2 .branch():按条件路由
场景:客服工单先做紧急度分级,再路由到不同处理步骤。
typescript
// src/mastra/workflows/ticket-workflow.ts(节选)
import { createStep, createWorkflow } from '@mastra/core/workflows';
import { z } from 'zod';
// 步骤一:分类打标
const classifyStep = createStep({
id: 'classify',
inputSchema: z.object({ text: z.string() }),
outputSchema: z.object({
text: z.string(),
urgent: z.boolean(), // 分类结果作为路由依据
}),
execute: async ({ inputData }) => {
const urgent = /宕机|无法支付|数据丢失/.test(inputData.text);
return { text: inputData.text, urgent };
},
});
const urgentHandler = createStep({
id: 'urgent-handler',
inputSchema: z.object({ text: z.string(), urgent: z.boolean() }),
outputSchema: z.object({ handled: z.boolean(), channel: z.string() }),
execute: async () => ({ handled: true, channel: 'pager-值班' }),
});
const normalQueue = createStep({
id: 'normal-queue',
inputSchema: z.object({ text: z.string(), urgent: z.boolean() }),
outputSchema: z.object({ handled: z.boolean(), channel: z.string() }),
execute: async () => ({ handled: true, channel: '普通队列' }),
});typescript
// 组装分支:每个分支是 [过滤条件, 目标步骤] 的元组
export const ticketWorkflow = createWorkflow({
id: 'ticket',
inputSchema: z.object({ text: z.string() }),
outputSchema: z.object({ handled: z.boolean(), channel: z.string() }),
})
.then(classifyStep)
.branch([
// 分支 1:urgent 为 true 时走紧急处理
[async ({ inputData }) => inputData.urgent === true, urgentHandler],
// 分支 2:否则进普通队列
[async ({ inputData }) => inputData.urgent === false, normalQueue],
])
.commit();每个分支的条件函数接收上下文并返回布尔值,命中的分支被执行,未命中被跳过。
8.3 .parallel():并行扇出
场景:内容发布前需要同时做"敏感词审查"和"SEO 评分",两者互不依赖,并行执行省一半耗时:
typescript
// 两个并行步骤:输入相同、输出各自独立
export const publishCheck = createWorkflow({
id: 'publish-check',
inputSchema: z.object({ article: z.string() }),
outputSchema: z.object({
review: z.object({ pass: z.boolean() }),
seo: z.object({ score: z.number() }),
}),
})
.parallel([
// 支路 1:敏感词审查
moderationStep,
// 支路 2:SEO 评分
seoScoreStep,
])
.commit();并行块的输出会按键分组聚合:每条支路的输出挂在各自的 step id 下,供下游读取——下游步骤的 inputSchema 需要声明这个嵌套结构。
8.4 汇合后聚合
typescript
// 汇合步骤:消费 parallel 的聚合输出
const mergeStep = createStep({
id: 'merge',
inputSchema: z.object({
// 键名 = 各支路步骤的 id
moderation: z.object({ pass: z.boolean() }),
seo: z.object({ score: z.number() }),
}),
outputSchema: z.object({ ok: z.boolean() }),
execute: async ({ inputData }) => {
const ok =
inputData.moderation.pass && inputData.seo.score >= 60;
return { ok }; // 双重检查都通过才允许发布
},
});
// 完整链路:并行 → 汇合判定
export const fullPipeline = createWorkflow({
id: 'full-pipeline',
inputSchema: z.object({ article: z.string() }),
outputSchema: z.object({ ok: z.boolean() }),
})
.parallel([moderationStep, seoScoreStep])
.then(mergeStep)
.commit();性能提示
把互不依赖的 IO 密集步骤(多路 API 调用、多个模型请求)放进 parallel,总耗时约等于最慢一路而非全部之和。
本章小结
- 三种控制流:
.then()串行、.branch()条件路由、.parallel()并行扇出; - branch 的每个分支是
[条件函数, 步骤]元组,条件基于上游 inputData 判断; - parallel 的输出按支路 step id 聚合成对象,下游 schema 要声明嵌套结构;
- 无依赖的 IO 步骤尽量并行,可显著压缩端到端延迟。
🧪 随堂测验
点击你认为正确的选项。答错时会展示正确答案与原因解析。
1. .branch() 的每个分支元素的结构是?
2. 两个互相独立的外部 API 调用步骤,最优的组织方式是?
3. .parallel([aStep, bStep]) 之后,下游步骤的 inputSchema 应如何声明?
4. branch 中所有条件函数都不命中会发生什么?
🛠️ 动手实践
- 把 ticketWorkflow 扩展成三分支:紧急 / 普通 / 垃圾广告(新增 spam 过滤条件)。
- 给 fullPipeline 的 parallel 再加一条支路
readabilityStep(可读性评分),更新 mergeStep 的 schema 与判定逻辑。 - 在 Studio 中分别触发 urgent=true 和 false 的工单,观察 Graph 视图中实际走过的路径高亮。