Skip to content

第 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 中所有条件函数都不命中会发生什么?

🛠️ 动手实践

  1. 把 ticketWorkflow 扩展成三分支:紧急 / 普通 / 垃圾广告(新增 spam 过滤条件)。
  2. 给 fullPipeline 的 parallel 再加一条支路 readabilityStep(可读性评分),更新 mergeStep 的 schema 与判定逻辑。
  3. 在 Studio 中分别触发 urgent=true 和 false 的工单,观察 Graph 视图中实际走过的路径高亮。

下一章:第 9 章 · Suspend & Resume 人机协同