跳到主要内容

Temporal Workflow

Temporal 是一个持久执行平台,用于编排长时运行、可容错的 Workflow。@mastra/temporal 包让你可以使用标准 Mastra API 编写 Workflow,并在 Temporal 集群上运行。

注意

@mastra/temporal 仍处于实验阶段,尚不适合生产使用。API 可能会在不同版本之间发生变化。当前状态请参阅包 README

Temporal 与 Mastra 的工作原理
Temporal 与 Mastra 的工作原理的直接链接

使用 createWorkflow()createStep() 编写的 Mastra Workflow 会映射到 Temporal 的 Workflow 和 Activity 模型。用于 Temporal Worker 的 MastraPlugin 会在 bundling 时编译 Mastra 入口文件:

  • 每个 createStep() handler 都会提取为 Temporal Activity。
  • 每个 createWorkflow() 都会重写为调用这些 Activity 的 Temporal Workflow。
  • 该插件会自动向 Worker 注册生成的 Activity 和 Workflow。

通过 mastra.getWorkflow(...).createRun().start(...) 开始运行时,Mastra Client 会将控制权交给 Temporal。随后,Temporal 在 Worker 上驱动持久执行、重试和状态持久化。

设置
设置的直接链接

安装所需包:

npm install @mastra/temporal@latest @temporalio/client @temporalio/worker @temporalio/envconfig

你还需要能够访问 Temporal 集群。本地开发时,可以使用 Docker 运行集群(参阅在本地运行)。

构建由 Temporal 支持的 Workflow
构建由 Temporal 支持的 Workflow的直接链接

本指南逐步介绍如何使用 Temporal 和 Mastra 创建 Workflow,并通过一个数值递增的计数器应用进行演示。

初始化 Temporal
初始化 Temporal的直接链接

初始化 Temporal 集成,以获取与 Mastra 兼容的 Workflow 辅助函数。createWorkflow()createStep() 函数会绑定到 Temporal Client 和任务队列。

src/mastra/temporal/index.ts
import { init } from '@mastra/temporal'
import { Client, Connection } from '@temporalio/client'
import { loadClientConnectConfig } from '@temporalio/envconfig'

const config = loadClientConnectConfig()
const connection = await Connection.connect(config.connectionOptions)
const client = new Client({ connection })

export const { createWorkflow, createStep } = init({
client,
taskQueue: 'mastra',
})

loadClientConnectConfig() 会读取标准 Temporal 环境变量,例如 TEMPORAL_ADDRESSTEMPORAL_NAMESPACE 和 mTLS 设置。完整列表请参阅 Temporal envconfig 文档

创建步骤
创建步骤的直接链接

定义组成 Workflow 的各个步骤。每个步骤都会成为 Temporal Activity。

src/mastra/workflows/index.ts
import { z } from 'zod'
import { createWorkflow, createStep } from '../temporal'

const incrementStep = createStep({
id: 'increment',
inputSchema: z.object({
value: z.number(),
}),
outputSchema: z.object({
value: z.number(),
}),
execute: async ({ inputData }) => {
return { value: inputData.value + 1 }
},
})

创建 Workflow
创建 Workflow的直接链接

将这些步骤组合成 Workflow。Workflow id 必须是静态字符串字面量,以便构建时 transformer 推导其 Temporal export 名称。

src/mastra/workflows/index.ts
const workflow = createWorkflow({
id: 'increment-workflow',
steps: [incrementStep],
inputSchema: z.object({
value: z.number(),
}),
outputSchema: z.object({
value: z.number(),
}),
}).then(incrementStep)

workflow.commit()

export { workflow as incrementWorkflow }

配置 Mastra 实例
配置 Mastra 实例的直接链接

向 Mastra 注册 Workflow。执行由 Temporal Worker 驱动。

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { PinoLogger } from '@mastra/loggers'
import { incrementWorkflow } from './workflows'

export const mastra = new Mastra({
workflows: { incrementWorkflow },
logger: new PinoLogger({ name: 'Mastra', level: 'info' }),
})

运行 Worker
运行 Worker的直接链接

Worker 是一个轮询 Temporal 任务队列的长时运行 Node.js 进程。安装 MastraPlugin,并将其 src 选项指向注册 Workflow 的 Mastra 入口文件。

src/mastra/worker.ts
import { MastraPlugin } from '@mastra/temporal/worker'
import { NativeConnection, Worker } from '@temporalio/worker'

const connection = await NativeConnection.connect({
address: 'localhost:7233',
})

const mastraPlugin = new MastraPlugin()

await mastraPlugin.prebuild({
entryFile: import.meta.resolve('./index.ts'),
})

const worker = await Worker.create({
connection,
namespace: 'default',
taskQueue: 'mastra',
plugins: [mastraPlugin],
})

await worker.run()

MastraPlugin 会将入口文件重写为仅含 Workflow 的 bundle,并将步骤 handler 接入为 Temporal Activity。无需手动向 Worker.create() 传递 activitiesworkflowsPath

运行 Workflow
运行 Workflow的直接链接

在本地运行
在本地运行的直接链接

  1. 启动本地 Temporal Server。最简单的选择是使用 temporalio/auto-setup Docker 镜像:

    docker run --rm -p 7233:7233 -p 8080:8080 temporalio/auto-setup:latest
  2. http://localhost:8080 打开 Temporal UI,检查 namespace、Workflow 和 Activity。

  3. 在新终端中运行以下命令启动 Worker:

    npx tsx src/mastra/worker.ts
  4. 从脚本或任何导入 Mastra 实例的进程中触发 Workflow 运行:

    scripts/run.ts
    import { mastra } from '../src/mastra'

    const run = await mastra.getWorkflow('incrementWorkflow').createRun()
    const result = await run.start({ inputData: { value: 5 } })

    console.log(result)
  5. 在 Temporal UI 的 Workflows 下监控执行,查看逐步 Activity 进度和重试历史。

在生产环境中运行
在生产环境中运行的直接链接

在生产环境中,请使用 Temporal Cloud 或自行托管的 Temporal 集群。通过 @temporalio/envconfig 读取的环境变量配置 Client 和 Worker 连接:

.env
TEMPORAL_ADDRESS=your-namespace.tmprl.cloud:7233
TEMPORAL_NAMESPACE=your-namespace
TEMPORAL_API_KEY=your-api-key

有关 mTLS 和 API Key 选项,请参阅 Temporal Cloud 连接文档

注意

Temporal Worker 必须作为长时运行的进程运行。请勿将其部署到 AWS Lambda 或 Vercel Function 等执行时间限制较短的 Serverless 平台。请使用容器、虚拟机,或 Fly.io、Railway、Kubernetes 等适合 Worker 的平台。

配置选项
配置选项的直接链接

taskQueue
taskqueue的直接链接

必填。标识 Worker 轮询的 Temporal 任务队列。必须将相同值传给 init()(供 Client 启动运行)和 Worker.create()(供 Worker 接收运行)。

startToCloseTimeout
starttoclosetimeout的直接链接

可选。设置单个 Activity(步骤)允许运行的最长时间;超过后,Temporal 会取消该 Activity 并应用重试策略。默认值为 1 minute

src/mastra/temporal/index.ts
export const { createWorkflow, createStep } = init({
client,
taskQueue: 'mastra',
startToCloseTimeout: '5 minutes',
})

限制与说明
限制与说明的直接链接

  • Workflow ID 必须是静态字符串字面量。构建时 transformer 会读取该字面量值,以推导 Temporal Workflow export 名称。
  • Activity 会根据 createStep() handler 自动生成。请勿在 Worker.create({ activities }) 中传入。