Snapshot
在 Mastra 中,Snapshot 是 Workflow 在特定时间点的完整执行状态的可序列化表示。Snapshot 会捕获从原位置恢复 Workflow 所需的全部信息,包括:
- Workflow 中每个步骤的当前状态
- 已完成步骤的输出
- Workflow 实际采用的执行路径
- 所有已挂起步骤及其元数据
- 每个步骤剩余的重试次数
- 恢复执行所需的其他上下文数据
每当 Workflow 挂起时,Mastra 都会自动创建和管理 Snapshot,并将其持久化到配置的 Storage 系统。
Snapshot 在挂起与恢复中的作用Snapshot 在挂起与恢复中的作用的直接链接
Snapshot 是实现 Mastra 挂起与恢复能力的关键机制。当 Workflow 步骤调用 await suspend() 时:
- Workflow 执行会在该确切位置暂停
- Workflow 当前状态会被捕获为 Snapshot
- Snapshot 会持久化到 Storage
- Workflow 步骤会标记为“已挂起”,状态为
'suspended' - 稍后在已挂起步骤上调用
resume()时,会检索 Snapshot - Workflow 执行会从原位置精确恢复
此机制为实现 Human-in-the-loop Workflow、处理速率限制、等待外部 Resource,以及实现可能需要长时间暂停的复杂分支 Workflow 提供了强大方式。
Snapshot 结构Snapshot 结构的直接链接
每个 Snapshot 都包含 runId、输入、步骤状态(success、suspended 等)、所有挂起和恢复载荷以及最终输出。因此,恢复执行时可以获得完整上下文。
{
"runId": "34904c14-e79e-4a12-9804-9655d4616c50",
"status": "success",
"value": {},
"context": {
"input": {
"value": 100,
"user": "Michael",
"requiredApprovers": ["manager", "finance"]
},
"approval-step": {
"payload": {
"value": 100,
"user": "Michael",
"requiredApprovers": ["manager", "finance"]
},
"startedAt": 1758027577955,
"status": "success",
"suspendPayload": {
"message": "Workflow suspended",
"requestedBy": "Michael",
"approvers": ["manager", "finance"]
},
"suspendedAt": 1758027578065,
"resumePayload": { "confirm": true, "approver": "manager" },
"resumedAt": 1758027578517,
"output": { "value": 100, "approved": true },
"endedAt": 1758027578634
}
},
"activePaths": [],
"serializedStepGraph": [
{
"type": "step",
"step": {
"id": "approval-step",
"description": "Accepts a value, waits for confirmation"
}
}
],
"suspendedPaths": {},
"waitingPaths": {},
"result": { "value": 100, "approved": true },
"requestContext": {},
"timestamp": 1758027578740
}
Snapshot 如何保存和检索Snapshot 如何保存和检索的直接链接
Snapshot 会保存到配置的 Storage 系统。默认使用 libSQL,但也可以改为配置 Upstash、PostgreSQL 或 OracleDB。每个 Snapshot 都保存在 workflow_snapshots 表中,并通过 Workflow 的 runId 标识。
了解更多:
保存 Snapshot保存 Snapshot的直接链接
Workflow 挂起时,Mastra 会通过以下步骤自动持久化 Workflow Snapshot:
- 步骤执行中的
suspend()函数触发 Snapshot 流程 WorkflowInstance.suspend()方法记录已挂起的状态机- 调用
persistWorkflowSnapshot()保存当前状态 - Snapshot 被序列化,并存入配置数据库的
workflow_snapshots表 - Storage 记录包含 Workflow 名称、Run ID 和序列化的 Snapshot
检索 Snapshot检索 Snapshot的直接链接
Workflow 恢复时,Mastra 会通过以下步骤检索持久化的 Snapshot:
- 使用特定步骤 ID 调用
resume()方法 - 使用
loadWorkflowSnapshot()从 Storage 加载 Snapshot - 解析 Snapshot 并准备恢复
- 使用 Snapshot 状态重新创建 Workflow 执行
- 恢复已挂起步骤并继续执行
const storage = mastra.getStorage()
const workflowStore = await storage?.getStore('workflows')
const snapshot = await workflowStore?.loadWorkflowSnapshot({
runId: '<run-id>',
workflowName: '<workflow-id>',
})
console.log(snapshot)
Snapshot 的 Storage 选项Snapshot 的 Storage 选项的直接链接
Snapshot 使用 Mastra 类上配置的 storage 实例持久化。该 Storage 层由注册到该实例的所有 Workflow 共享。Mastra 支持多种 Storage 选项,以适应不同环境。
import { Mastra } from '@mastra/core'
import { LibSQLStore } from '@mastra/libsql'
import { approvalWorkflow } from './workflows'
export const mastra = new Mastra({
storage: new LibSQLStore({
id: 'mastra-storage',
url: ':memory:',
}),
workflows: { approvalWorkflow },
})
- libSQL Storage
- PostgreSQL Storage
- OracleDB Storage
- MongoDB Storage
- Upstash Storage
- Cloudflare D1
- DynamoDB
- 更多 Storage Provider
最佳实践最佳实践的直接链接
- 确保可序列化:需要包含在 Snapshot 中的所有数据都必须可序列化(能够转换为 JSON)。
- 尽量缩小 Snapshot:避免直接在 Workflow 上下文中存储大型数据对象。应存储对它们的引用(例如 ID),并在需要时检索数据。
- 谨慎处理恢复上下文:恢复 Workflow 时,请仔细考虑要提供的上下文,因为它会与现有 Snapshot 数据合并。
- 设置适当的监控:为已挂起 Workflow(尤其是长期运行的 Workflow)实施监控,并确保它们能够正确恢复。
- 考虑 Storage 扩展能力:对于包含大量已挂起 Workflow 的应用,请确保 Storage 方案得到适当扩展。
自定义 Snapshot 元数据自定义 Snapshot 元数据的直接链接
可以通过定义 suspendSchema,在挂起 Workflow 时附加自定义元数据。这些元数据会存储在 Snapshot 中,并在 Workflow 恢复时可用。
import { createWorkflow, createStep } from '@mastra/core/workflows'
import { z } from 'zod'
const approvalStep = createStep({
id: 'approval-step',
description: 'Accepts a value, waits for confirmation',
inputSchema: z.object({
value: z.number(),
user: z.string(),
requiredApprovers: z.array(z.string()),
}),
suspendSchema: z.object({
message: z.string(),
requestedBy: z.string(),
approvers: z.array(z.string()),
}),
resumeSchema: z.object({
confirm: z.boolean(),
approver: z.string(),
}),
outputSchema: z.object({
value: z.number(),
approved: z.boolean(),
}),
execute: async ({ inputData, resumeData, suspend }) => {
const { value, user, requiredApprovers } = inputData
const { confirm } = resumeData ?? {}
if (!confirm) {
return await suspend({
message: 'Workflow suspended',
requestedBy: user,
approvers: [...requiredApprovers],
})
}
return {
value,
approved: confirm,
}
},
})
提供恢复数据提供恢复数据的直接链接
使用 resumeData 在恢复已挂起步骤时传入结构化输入。它必须与步骤的 resumeSchema 匹配。
const workflow = mastra.getWorkflow('approvalWorkflow')
const run = await workflow.createRun()
const result = await run.start({
inputData: {
value: 100,
user: 'Michael',
requiredApprovers: ['manager', 'finance'],
},
})
if (result.status === 'suspended') {
const resumedResult = await run.resume({
step: 'approval-step',
resumeData: {
confirm: true,
approver: 'manager',
},
})
}