基于 Cloudflare D1 构建 Mastra 存储层:@mastra/cloudflare-d1 完整实战与版本演进解读
【免费下载链接】mastraMastra is the modern TypeScript framework for AI-powered applications and agents.项目地址: https://gitcode.com/GitHub_Trending/ma/mastra
@mastra/cloudflare-d1是 Mastra 面向 Cloudflare D1(基于 SQLite 的全球分布式数据库)的官方存储适配器,为线程(Threads)、消息(Messages)、工作流(Workflows)、评分(Scores)与后台任务(Background Tasks)提供持久化能力。本文以该包的官方变更记录(CHANGELOG.md)为骨架,结合仓库源码(src/storage/index.ts、db/index.ts)深入讲解三种接入方式、域存储架构、D1 特性相关的底层实现原理,以及从 1.0 到 1.3 的重要 API 演进与迁移路径,读完即可在 Workers/Pages 场景中正确选型并落地 Mastra 持久化。
一、包定位与能力全景
1.1 它解决什么问题
Mastra 中,Agent 的对话历史(Thread/Memory)、工作流快照、评估分数与链路追踪都需要落库。@mastra/cloudflare-d1让开发者可以在 Cloudflare Workers 生态内直接使用 D1 作为唯一的持久化后端,无需额外部署独立数据库,天然享有 D1 的全球读副本、按量计费与 SQLite 兼容性。包的 package.json 声明了它的依赖画像:
- 运行时依赖仅
cloudflare(^5.2.0),用于走 REST API 时构造官方客户端; peerDependencies为@mastra/core(>=1.54.0-0 <2.0.0-0)与@cloudflare/workers-types(^4.20240919.0);- 引擎要求 Node.js >= 22.13.0。
从源码结构看,包以域(domain)驱动组织:src/storage/domains/下分为memory、workflows、scores、background-tasks四个域实现,每个域继承@mastra/core提供的基类,最终由D1Store组合成一个复合存储(Composite Store)。
1.2 版本脉络一览
| 版本 | 类型 | 核心变化 |
|---|---|---|
| 0.12.x | Patch | ESM 声明修复、getMessagesById引入、快照加载最近一条、peerDeps 维护 |
| 0.13.x | Patch | 新增按 trace/span id 拉取分数、Workflow 快照携带 resourceId、空 threadId 校验 |
| 1.0.0 | Major | 存储架构重构为域存储(getStore()模式)、统一分页page/perPage、移除旧 API |
| 1.1.x | Patch | 评分增加batchId/datasetId/datasetItemId与多租户字段 |
| 1.2.x | Minor | 消息历史支持精确元数据过滤(memory.recallfilter.metadata) |
| 1.3.x | Patch/Minor | 内存读取错误显式抛MastraError、部分线程更新能力、后台任务 CAS 语义、@mastra/core 版本对齐 |
下文按“接入配置 → 域存储用法 → 平台特性源码解析 → 重要 API 演进 → 升级迁移”的顺序展开。
二、三种接入方式与完整配置参数
D1Store支持三种初始化方式,源码 src/storage/index.ts 通过配置联合类型D1StoreConfig = D1Config | D1WorkersConfig | D1ClientConfig统一处理,构造函数根据传入字段自动选择执行路径。
2.1 Workers Binding 方式(推荐)
在wrangler.toml中绑定 D1 数据库后,将env.DB直接注入:
import { D1Store } from '@mastra/cloudflare-d1'; const store = new D1Store({ binding: env.DB, // D1Database binding,来自 Worker 环境 tablePrefix: 'mastra_', // 可选,表名前缀 }); export default { async fetch(request, env) { const store = new D1Store({ binding: env.DB }); // ... }, };2.2 REST API 方式
适用于非 Workers 运行环境(如本地 Node 脚本、CI)通过 Cloudflare REST API 访问 D1:
const store = new D1Store({ id: 'my-d1-store', accountId: process.env.CF_ACCOUNT_ID, apiToken: process.env.CF_API_TOKEN, databaseId: process.env.CF_D1_DATABASE_ID, tablePrefix: 'mastra_', // 可选 });源码中该分支会用new Cloudflare({ apiToken })构造官方客户端,并把sql与params透传给cfClient.d1.database.query(databaseId, ...)。
2.3 预配置 Client 方式
若需要自定义连接行为(超时、拦截器等),可传入已配置好的 client:
const store = new D1Store({ id: 'my-d1-store', client: { query: async ({ sql, params }) => { // 自定义实现或包装官方客户端 return cfClient.d1.database.query(databaseId, { account_id, sql, params }); }, }, });2.4 通用配置项
| 配置项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
id | string | 必填 | 存储实例 ID,用于日志与错误追踪 |
tablePrefix | string | '' | 所有表名的前缀,只允许字母、数字、下划线(构造时用正则/^[a-zA-Z0-9_]*$/校验,非法值直接抛错) |
disableInit | boolean | false | 为 true 时关闭运行时自动建表/迁移,适合 CI/CD 中以高权限先执行storage.init()、运行时再以只读权限启动 |
disableInit的完整语义在源码注释中有明确说明(db/index.ts):开启后存储不会在首次使用时自动创建/修改表结构,必须在部署脚本中显式调用await storage.init()。
三、域存储架构:getStore() 模式
CHANGELOG 1.0.0 记录了一次重大架构重构:将MastraStorage基类上的透传方法改为“域存储(Domain Store)”模式,所有适配器统一通过getStore('domainName')访问具体能力,D1Store继承MastraCompositeStore并暴露stores属性(源码见 src/storage/index.ts)。
// Before(已移除) const thread = await storage.getThreadById({ threadId }); // After const memory = await storage.getStore('memory'); const thread = await memory?.getThreadById({ threadId }); const workflows = await storage.getStore('workflows'); await workflows?.persistWorkflowSnapshot({ workflowName, runId, snapshot }); const scores = await storage.getStore('scores'); await scores?.saveScore({ ...score });D1Store内置四个域:memory、workflows、scores、backgroundTasks。每个域实现都接受统一格式的D1DomainConfig(client / binding / REST 三种形态之一,由 resolveD1Config 归一化后构造共享的D1DB实例),从而让多个域复用同一条连接与同一套 SQL 构建逻辑。
3.1 各域能力对照
| 域 | 基类 | 主要方法 | 依赖表 |
|---|---|---|---|
| memory | MemoryStorage | saveThread、getThreadById、updateThread、deleteThread、listThreads、saveMessages、listMessages、listMessagesById、updateMessages、saveResource | threads、messages、resources |
| workflows | WorkflowsStorage | persistWorkflowSnapshot、loadWorkflowSnapshot、listWorkflowRuns、getWorkflowRunById | workflow_snapshot |
| scores | ScoresStorage | saveScore、getScoreById、listScoresByScorerId等 | scorers |
| backgroundTasks | 后台任务域 | 任务状态存储与查询 | background_tasks 相关表 |
dangerouslyClearAll()为每个域都实现了清表逻辑(delete全部行而非 DROP),供测试清理使用。
四、D1 平台特性下的源码级实现
D1 本质是 SQLite,因此D1DB层(src/storage/db/index.ts)做了大量 SQLite 特性适配。理解这些实现,有助于预判 D1 后端的边界行为。
4.1 SQLite 类型映射
getSqlType()将通用 StorageColumn 类型映射为 SQLite 类型(db/index.ts):
| 通用类型 | SQLite 映射 | 原因 |
|---|---|---|
bigint | INTEGER | SQLite 整数不分大小 |
jsonb | TEXT | JSON 以文本存储 |
boolean | INTEGER | 以 0/1 存储 |
| 其他 | 透传getSqlType | 保持默认映射 |
4.2 弹性列处理(前向兼容)
CHANGELOG 1.0.2 提到“未知列静默丢弃”。实现上,D1DB会通过PRAGMA table_info(表名)读取真实列集合并缓存(tableColumnsCache),insert/update/batchInsert在拼 SQL 前用filterRecordToKnownColumns过滤掉表中不存在的字段(db/index.ts)。这意味着:
// 即使 record 里多了未来版本才有的 futureField,也不会抛 SQL 错误 await db.insert({ tableName, record: { id: '1', title: 'Hello', futureField: 'value' } });这保证了“新版本 domain 包先写新字段、旧版本存储表尚未迁移”的滚动升级场景不崩溃。DDL 执行(alterTable/dropTable/createTable)后缓存会被主动失效,避免读到过期 schema。
4.3 复合 SELECT 限制与 UNION ALL 批处理
CHANGELOG 1.0.3 修复了“语义召回 +perPage=0+ 多个 include 目标时返回空结果”的问题。根因是 D1 对SQLITE_LIMIT_COMPOUND_SELECT的限制远低于 SQLite 默认的 500(实测约 10 项即失败)。_getIncludedMessages(memory/index.ts)因此:
- 先批量拉取目标消息(
WHERE id IN (...)),消除按子查询关联的写法; - 以
MAX_UNION_TERMS = 5为批大小,把每个 include 的“前文窗口 + 后文窗口”子查询用UNION ALL分批拼接执行; - 最终按
createdAt ASC, id ASC统一排序并去重。
const messages = await memory.recall({ threadId: 'thread-1', perPage: 0, // 语义召回:只要 include 命中的上下文 include: [{ id: 'msg-3', withPreviousMessages: 2, withNextMessages: 1 }], });4.4 JSON 元数据过滤
D1 的元数据以 JSON 文本存于列中,listThreads/listMessages使用 SQLite 的json_valid、json_type、json_extract函数实现类型安全的过滤(memory/index.ts),并区分null、字符串、数字、布尔四种值类型:
const messages = await memory.recall({ threadId: 'thread-1', filter: { metadata: { status: 'done', // 多字段之间为 AND 语义 priority: 'high', }, }, }); // 支持的值类型:字符串、有限数字、布尔值、null安全加固:元数据 key 与分页参数都经过校验(validateMetadataKeys、validatePaginationInput),防止 SQL 注入与整数溢出攻击(CHANGELOG 1.0.0listThreads条目)。
4.5 并发更新的明确边界
D1 无法提供原子的读-改-写语义(batch 虽原子但必须一次性提交全部语句),因此工作流域的updateWorkflowResults与updateWorkflowState直接抛出 not-implemented 错误(workflows/index.ts),supportsConcurrentUpdates()返回false。同样,CHANGELOG 1.3.1 明确指出:后台任务的原子条件状态更新(CAS)语义不再暴露给 Cloudflare KV 与 ClickHouse,因为它们无法保证比较-交换的一致性。
五、错误语义与可观测性
5.1 统一错误 ID
自 1.0.0 起,所有存储与向量库统一使用createStorageErrorId,错误 ID 形如MASTRA_STORAGE_{STORE}_{OPERATION}_{STATUS}(例如MASTRA_STORAGE_CLOUDFLARE_D1_LIST_MESSAGES_FAILED)。D1 的每个操作(CREATE_TABLE、INSERT、BATCH_UPSERT、LOAD、ALTER_TABLE等)失败时都会包装成带domain: ErrorDomain.STORAGE、category: ErrorCategory.THIRD_PARTY | USER | SYSTEM的MastraError,方便日志聚合与告警定位。
5.2 内存读取错误不再静默吞掉
CHANGELOG 1.3.0(Minor)是一个重要的行为变更:分页内存读取(listThreads、listMessages、listMessagesByResourceId、listMessagesById)此前在后台失败(表锁、连接断开)时会记录日志并返回{ threads: [], total: 0, hasMore: false }之类的空载荷,导致 Agent 把瞬时故障误判为“没有历史”,可能覆盖真实状态。1.3.0 起这些方法改抛MastraError(校验类 USER 错误与真正的空结果行为不变):
try { const { threads } = await storage.listThreads({ resourceId }); // ...使用 threads } catch (error) { // 真实的后端故障:决定是重试、透传还是降级 // 空列表只表示“确实没有线程”,不再混入故障 }六、重要 API 演进与迁移路径
CHANGELOG 记录了从 0.x 到 1.3 的多次破坏性变更,升级时请重点核对以下几处。
6.1 分页:offset/limit → page/perPage(1.0.0)
所有存储与内存分页 API 统一改为 0 起始的page与perPage,且perPage支持false表示“不分页取全部”:
// Before await memory.listThreadsByResourceId({ resourceId: 'user-123', offset: 20, limit: 10 }); // After await memory.listThreadsByResourceId({ resourceId: 'user-123', page: 2, perPage: 10 }); // 一次性取全部消息 await storage.listMessages({ threadId: 'thread-1', page: 0, perPage: false });新增校验:负page抛错、perPage的负值/0/false 边界统一处理。
6.2 消息读取 API 收敛(1.0.0)
- 移除
getMessagesPaginated()与getMessages(),统一使用listMessages()(默认按createdAt升序,可用orderBy: { field: 'createdAt', direction: 'DESC' }倒序); listMessages()要求非空、非纯空白的threadId(否则抛错而非返回空);getMessagesById({ messageIds, format })→listMessagesById({ messageIds }),只返回 V2 格式消息;- 类型
StorageGetMessagesArg→StorageListMessagesInput; - 客户端 SDK
client.getThreadMessages()→client.listThreadMessages()。
6.3 复合存储与命名演进
- 1.0.0-beta 引入
StorageDomain基类与InMemoryDB,域可独立使用(传配置自建 client,或传已有 client 共享连接); - 1.0.0 支持
MastraStorage组合不同适配器的域(如 memory 用 libsql、workflows 用 PG); - 1.0.0-beta.12 / 1.0.0 将
MastraStorage更名为MastraCompositeStore(旧名保留为弃用别名):
// Before import { MastraStorage } from '@mastra/core/storage'; // After import { MastraCompositeStore } from '@mastra/core/storage';6.4 部分线程更新与标题保护(1.3.0)
updateThread的title与metadata变为相互独立可选:只更新其中一个时,另一个列保持不变,修复了“轮次内新生成的标题被陈旧值覆盖”(#21041);- 存储适配器声明
supportsPartialThreadUpdate = true(源码见 memory/index.ts),旧适配器则由@mastra/memory回填既有标题,保证混合版本部署兼容; - 修复了无标题线程在观察性记忆缓冲期间写入 null 标题触发 not-null 约束崩溃的问题(#21257)。
6.5 评分的溯源与多租户字段(1.1.1)
scoreTrace()与saveScore()支持可选的batchId、datasetId、datasetItemId(把一次评分批次与数据集条目关联),以及organizationId、projectId(多租户隔离,projectId与表示记忆资源的resourceId语义分离):
await scoreTrace({ storage, scorer, target: { traceId }, batchId: 'baseline-batch-1', datasetId, datasetItemId, }); await storage.saveScore({ ...score, organizationId: 'org-a', projectId: 'proj-1' }); const result = await storage.listScoresByScorerId({ scorerId, filters: { organizationId: 'org-a', projectId: 'proj-1' }, });D1 的 scorers 表通过alterTable附加迁移新增这些列,并使用transformScoreRow统一行变换(D1 偏好createdAtZ/updatedAtZ时间戳字段,见 scores/index.ts)。
6.6 其他值得注意的变更
- 1.0.7:针对 2026-06-17 "easy-day-js" 供应链事件的补丁发布,重新发布干净版本并将
latestdist-tag 前移; - 1.2.0:修复多线程消息查询中 include 消息与分页元数据返回不正确的问题(#20303);
- 1.3.2:修正最低支持的
@mastra/core版本以匹配实际使用的 API,并移除发布包中的 CHANGELOG.md 以减小包体积; - 1.0.x 早期:为 SQL 域补齐
BackgroundTasksStorage域实现(#15307),并跟踪suspendedAt/suspendPayload(自动alterTable迁移); - 1.0.6:
getThreadById尊重可选resourceId——当线程属于其他资源时返回null。
七、测试与验证方式
仓库为 D1 适配器提供了较为完整的测试矩阵(均在stores/cloudflare-d1/src/storage/下):
sql-builder.test.ts:SQL 构建器单元测试;binding-api.test.ts/rest-api.test.ts:分别覆盖 Workers Binding 与 REST API 两条执行路径;memory-error-propagation.test.ts:验证 1.3.0 错误传播语义(列表读取失败抛MastraError);- 测试工具
test-utils.ts与 devDependencies 中的miniflare(^4.20260714.0)用于在本地模拟 Workers/D1 运行时,@internal/storage-test-utils提供跨适配器共享的存储断言。
本地运行测试:pnpm --filter @mastra/cloudflare-d1 test(包脚本定义见 package.json)。
八、总结与选型建议
- 首选 Workers Binding 接入:在 Worker/Pages Functions 内运行时免配置、延迟最低;REST 方式适合 CI、本地脚本与外部服务。
- 确认并发需求:若工作流依赖原子读-改-写(
updateWorkflowResults/updateWorkflowState)或需要后台任务 CAS 语义,D1 后端不适用,应选择具备事务/比较-交换能力的存储(如 PostgreSQL 系适配器)。 - 升级注意破坏性 API:从 0.x/1.0 早期升级时,重点迁移
page/perPage分页、listMessages*命名、getStore()域访问与MastraCompositeStore命名。 - 善用错误语义:1.3.0 起直接用
try/catch区分“后端故障”与“空结果”,避免 Agent 在瞬时故障时误清空记忆。
相关资源:README(安装与最小用法)、完整变更记录、核心实现 D1Store 与 D1DB。
【免费下载链接】mastraMastra is the modern TypeScript framework for AI-powered applications and agents.项目地址: https://gitcode.com/GitHub_Trending/ma/mastra
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考