Strapi Data Transfer 深入解析:Source Provider(ISourceProvider)接口与阶段流实现
本文以 Strapi 数据转移(data-transfer)实验特性中的 Source Provider 官方文档为主体,结合 ISourceProvider 接口定义、TransferEngine 引擎实现 与内置 source provider 源码,完整讲解 source provider 的接口结构、五个转移阶段各自产出的流式数据类型、引擎消费读流的调用链,以及编写自定义 source provider 的关键约束。读完后可掌握:如何实现一个符合 Strapi 数据转移规范的 source provider、每个 create{_stage}ReadStream() 方法的返回值语义,以及阶段流必须自行关闭的原因。
一、Source Provider 在数据转移中的定位
Strapi 的 data-transfer 包将一次数据迁移抽象为"源 → 引擎 → 目标"的流式管道:source provider 为每个转移阶段(stage)提供一个 Readable 读流,由引擎把读流、转换流(transform)、进度追踪流(tracker)与目标写流串联成 pipeline 逐阶段执行。这一分工在 Providers 概览文档 中有明确定义:
Source providers provide read streams for each stage in the transfer.
Source provider 是 provider 三元组中负责"取数"的一端,Strapi 内置了三种 source provider:
- Local Strapi:通过本地 Strapi 项目的数据库连接读取数据,见 local-source 实现;
- Strapi file:从一个标准化的 Strapi 转移文件包中读取,见 file source 实现;
- Directory / Remote Strapi:分别对应 directory source 与远程实例封装。
与 destination provider 提供 create{_stage}WriteStream() 写流相对,source provider 只需提供对应的读流方法;概览文档同时指出:不必同时实现 source 与 destination,只需实现自己场景需要的部分。
二、ISourceProvider 接口结构(文档核心要求)
官方文档要求:一个 source provider 必须实现 ISourceProvider 接口,该接口位于 packages/core/data-transfer/types/providers.d.ts(这是构建产物路径,对应的源码定义在 src/types/providers.ts)。接口的完整定义如下(摘自 providers.ts#L14-L48):
export type ProviderType = 'source' | 'destination';
export interface IProvider {
type: ProviderType;
name: string; // a unique name for this provider
results?: IProviderTransferResults;
/**
* bootstrap() is called during transfer engine bootstrap
* 用于初始化:建立数据库连接、打开文件、校验授权等
*/
bootstrap?(diagnostics?: IDiagnosticReporter): MaybePromise<void>;
close?(): MaybePromise<void>; // called during transfer engine close
getMetadata(): MaybePromise<IMetadata | null>; // 返回元数据,用于版本校验
getSchemas?(): MaybePromise<Record<string, Struct.Schema> | null>; // 返回 schema,用于一致性校验
beforeTransfer?(): MaybePromise<void>; // 在转移阶段执行前被调用
}
export interface ISourceProvider extends IProvider {
results?: ISourceProviderTransferResults;
/**
* Optional totals for a stage. Called by the engine after the stage read
* stream is created and before `stage::start`.
* 用于 CLI 进度展示(剩余字节/条数、ETA)。未知时可省略或返回 null。
*/
getStageTotals?(stage: TransferStage): MaybePromise<StageTotalsEstimate | null | undefined>;
createEntitiesReadStream?(): MaybePromise<Readable>;
createLinksReadStream?(): MaybePromise<Readable>;
createAssetsReadStream?(): MaybePromise<Readable>;
createConfigurationReadStream?(): MaybePromise<Readable>;
createSchemasReadStream?(): MaybePromise<Readable>;
}
从接口定义可以归纳出 source provider 的方法分三类:
| 分类 | 方法 | 是否必需 | 作用 |
|---|---|---|---|
| 身份标识 | type、name |
必需(类型字段) | type 固定为 'source';name 是 provider 的唯一名称 |
| 生命周期 | bootstrap(diagnostics?)、close()、beforeTransfer() |
可选 | 引擎在初始化、收尾、阶段执行前分别调用 |
| 元信息 | getMetadata()、getSchemas()、getStageTotals(stage) |
getMetadata 必需,其余可选 |
提供版本信息、schema 与阶段总量预估 |
| 阶段读流 | createEntitiesReadStream()、createLinksReadStream()、createAssetsReadStream()、createConfigurationReadStream()、createSchemasReadStream() |
均为可选方法 | 每个阶段一个 Readable 流,逐个 write(entity) 输出数据 |
需要注意两点源码级事实:
- 引擎在构造时会对两个 provider 做结构校验,见 engine/index.ts#L152-L163 中的
validateProvider('source', sourceProvider),因此实现必须满足接口契约(尤其是type === 'source')。 - 五个
create*ReadStream方法在接口上都带?(可选),引擎内部调用时也全部使用可选链(如this.sourceProvider.createEntitiesReadStream?.())。这意味着一个 source provider 可以只暴露部分阶段的读流,未实现的阶段会被引擎跳过(详见第四节的shouldSkipStage与 skip 分支)。
三、五个阶段读流与数据类型
文档的核心表述是:source provider "为每个阶段提供一组 create{_stage}ReadStream() 方法……对每个 entity、link(relation)、asset(file)、configuration entity 或 content type schema,根据所处阶段执行 stream.write(entity)"。各阶段对应的数据类型定义在 types/utils.ts 的 TransferStageTypeMap 与 types/common-entities.ts:
| 阶段 stage | 读流方法 | 流内每条数据的类型 | 说明 |
|---|---|---|---|
schemas |
createSchemasReadStream() |
Struct.Schema |
内容类型/组件的 schema 定义 |
entities |
createEntitiesReadStream() |
IEntity |
不含关联的内容数据:{ type, id, data } |
links |
createLinksReadStream() |
ILink(IBasicLink | IMorphLink | ICircularLink) |
实体间关系,kind 为 relation.basic / relation.morph / relation.circular |
assets |
createAssetsReadStream() |
IAsset |
媒体文件:{ filename, filepath, stream, stats, metadata },其中 stream 是二进制 Readable |
configuration |
createConfigurationReadStream() |
IConfiguration(core-store / webhook) |
项目配置、核心存储、webhook |
其中几个关键类型(摘自 common-entities.ts):
export interface IEntity<T extends UID.ContentType = UID.ContentType> {
type: T; // 父类型 UID(content-type、component 等)
id: number; // 实体引用 ID
data: Data.Entity<T>; // 实体属性值
}
export interface IAsset {
filename: string;
filepath: string;
stream: Readable; // 文件二进制流
stats: IAssetStats; // { size: number }
metadata: IFile; // 媒体元数据(hash、mime、size、formats 等)
buffer?: Buffer;
}
阶段执行顺序
引擎的实际执行顺序定义在 transfer() 方法中(engine/index.ts#L819-L868):
bootstrap() → init()(解析双方 metadata)→ integrityCheck()(版本 + schema 一致性)
→ beforeTransfer()
→ transferSchemas() → transferEntities() → transferAssets() → transferLinks() → transferConfiguration()
→ close()
这与 Stream Lifecycle 文档 中描述的五阶段流程一致:先转 schema,再转不含关联的实体,随后是媒体文件、实体间关系,最后才是项目配置。另有一个冻结的阶段常量数组 TRANSFER_STAGES(engine/index.ts#L51-L57):['entities', 'links', 'assets', 'schemas', 'configuration'],用于阶段遍历场景。
流关闭是 source provider 的硬性责任
官方文档特别强调:"当某阶段的流发送完所有数据后,该流必须先被关闭,转移引擎才会继续下一个阶段。" 这一点可以从引擎的 #transferStage 实现得到印证(engine/index.ts#L590-L680):
const controller = new AbortController();
this.#currentStreamController = controller;
await pipeline(streams, { signal }); // streams = [source, transform?, tracker?, destination]
引擎使用 Node.js stream/promises 的 pipeline 串联 source → transform → tracker → destination。pipeline 会在上游流自然结束(end/关闭)后依次销毁下游流,并在整个链路完成后才 resolve;transfer() 是 await 每一个 transferXxx() 串行的,因此 source 流若不 end/close,当前阶段将一直挂起,后续阶段永远不会开始。实现方在读完全部数据后应调用 stream.end()(或在异步生成器耗尽后结束流),这是与"引擎负责 close"最容易混淆的地方:引擎只负责销毁(destroy)与错误传播,正常结束(end)由 source provider 自己触发。
pipeline 同时带来两个工程收益:任一环节出错时自动销毁整条链(引擎 catch 后还会 destination.destroy(error) 并上报 stage::error 事件),以及原生的背压(backpressure)传递——source 端写入过快时会被 transform/tracker 反压暂停,这正是 Stream Lifecycle 文档 中 "Data flows in chunks through the pipeline / Each stream can process data asynchronously" 的底层机制。
四、引擎如何调用 Source Provider 的各个钩子
结合 engine/index.ts 源码,可以把引擎对 source provider 的完整调用链梳理为:
- 构造期:
validateProvider('source', ...)结构校验(L155); bootstrap():以Promise.allSettled并发执行 source 与 destination 的bootstrap(this.diagnostics),source 侧可在此建立数据库连接、打开 Strapi 文件包、校验授权;任一侧 reject 会触发panic(L705-L716);init()→#resolveProviderResource():调用sourceProvider.getMetadata()取得IMetadata(结构为{ strapi?: { version }, createdAt? },见 common-entities.ts#L4-L10),并与 destination 的 metadata 一起用于#assertStrapiVersionIntegrity的版本策略校验(versionStrategy可选ignore/patch/minor/major,默认ignore,见 engine/index.ts#L440-L488)。source 若不提供版本元数据,版本一致性检查会被跳过;integrityCheck():两侧getSchemas()提供的 schema 集合经#assertSchemasMatching做策略化 diff(schemaStrategy默认strict),source 的 schema 输出质量直接影响迁移能否通过校验;beforeTransfer():先 source 后 destination 依次执行,异常可被已注册的 error handler 接管;- 每个阶段:先调
create{Stage}ReadStream()取读流;对assets阶段,引擎还会在stage::start前调用getStageTotals('assets'),把返回的{ totalBytes?, totalCount? }合并进进度数据(#mergeSourceStageTotals,L1017-L1039),供 CLI 显示总量与 ETA。接口注释明确说明这是可选能力:"Omit or return null when unknown (older remotes, file providers)"; - 阶段执行:
#transferStage中依次发出stage::start→ 每写入一条数据stage::progress→stage::finish(或stage::error/stage::skip)事件,事件挂在引擎的progress.stream(PassThrough)上;整体迁移另有transfer::init/start/finish/error事件; - 收尾:
close()并发调用双方close?.();若迁移中途抛错,引擎会调用 destination 的rollback并抛出原错误,source 侧的读流由pipeline的错误路径负责销毁。
另外,shouldSkipStage(L563-L588)基于 TransferGroupPresets(content = links+entities、files = assets、config = configuration)结合引擎选项 only / exclude 决定阶段是否执行;被跳过的阶段若读流已创建,引擎会等待其 close 后直接 destroy 并发出 stage::skip(L611-L636)。这解释了为什么接口里读流方法全部可选——source provider 只实现某几个阶段时,其余阶段会自然走 skip 分支。
五、编写一个 Source Provider:最小可用实现
综合接口契约与引擎行为,一个最小 source provider 需要:type: 'source' + 唯一 name + 同步返回非空的 getMetadata()(引擎将其用于版本比对,返回 null 会跳过校验)+ 至少一个阶段的读流方法。示例(针对 entities 阶段,数据可来自任意来源,此处仅示意契约):
import { Readable } from 'stream';
import type { ISourceProvider, IEntity, IMetadata } from '@strapi/core-data-transfer';
class MySourceProvider implements ISourceProvider {
type = 'source' as const;
name = 'my-custom-source';
// 引擎在 integrityCheck 中读取:strapi.version 参与版本策略校验
async getMetadata(): Promise<IMetadata> {
return { strapi: { version: process.env.MY_STRAPI_VERSION ?? '5.0.0' } };
}
// 引擎在 bootstrap 阶段调用:建立连接/打开资源
async bootstrap() {
// 例如:打开一个数据目录、连接数据库
}
// entities 阶段:逐条 write(IEntity),全部写完后 end() 关闭流
async createEntitiesReadStream(): Promise<Readable> {
const rows = await this.loadEntities(); // 任意取数逻辑
return Readable.from(
rows.map((row) => ({
type: row.uid, // 内容类型 UID
id: row.id, // 实体引用 ID
data: row.attributes, // 属性值(关联字段留待 links 阶段)
})) as AsyncIterable<IEntity>,
{ objectMode: true }
);
}
// links 阶段:写 ILink,kind 取 'relation.basic' | 'relation.morph' | 'relation.circular'
async createLinksReadStream(): Promise<Readable> {
return Readable.from(await this.loadLinks(), { objectMode: true });
}
// 无需实现的阶段(assets/schemas/configuration)直接省略,
// 引擎会走 shouldSkipStage/#transferStage 的 skip 分支
}
实现时需要注意的约束均来自源码事实:
- 流必须自我结束:如第三节所述,
Readable.from耗尽迭代器后会自动end(),手工write时则要在末尾显式end(); - assets 阶段的数据携带内层流:
IAsset.stream本身是 Readable,引擎的#progressTrackerChunks会把每条 asset 的stream包一层计数 Transform 并替换asset.stream(engine/index.ts#L338-L393),因此 source 侧应保证asset.stream可被单一消费者 pipe,且stats.size与metadata尽量完整(destination 端依赖它们做落盘与校验); results字段是"逃生舱":接口注释写明 results 是 "optional object for tracking any data needed from outside the engine",迁移完成后transfer()会返回{ source: sourceProvider.results, destination: ..., engine: progress },source provider 可把统计信息(如各类型实体计数)挂在这里暴露给调用方;- 诊断上报:
bootstrap(diagnostics)收到的IDiagnosticReporter可用于上报非致命信息,与引擎自身的reportError/reportWarning/reportInfo汇入同一诊断栈。
六、内置实现与测试佐证
仓库内有三套可直接参考的 source provider 实现,覆盖了三种典型取数场景:
| 实现 | 路径 | 特点 |
|---|---|---|
| Local Strapi source | local-source/index.ts | 从运行中 Strapi 的数据库读取,按 entities / links / assets / configuration 分文件实现各阶段读流 |
| Strapi file source | file/providers/source/index.ts | 从标准化 Strapi 文件包读取,实现了 getStageTotals(用于总量预估) |
| Directory source | directory/providers/source/index.ts | 面向目录形态的读流封装 |
对应的测试用例验证了文档所描述的"逐条 stream.write(entity) + 阶段结束"行为:entities.test.ts、links.test.ts、assets.test.ts、configuration.test.ts 分别断言各阶段读流输出的实体结构,engine.test.ts 则覆盖了引擎串联各阶段的整体流程。若要确认某个阶段读流的字段细节,直接阅读这些测试是最快的方式。
七、注意事项与适用边界
- 该文档带有
experimental标签(见原文档 front-matter),data-transfer 整体属于实验特性,接口可能随版本调整; - Providers 概览文档 明确说明:目前所有 data-transfer provider 仅处理本地媒体资产(
/upload文件夹),provider 媒体还在开发中,一切与资产转移相关的处理(含 Strapi 文件结构、restore 策略与资产回滚)都被视为unstable,近期可能变化; - source provider 只负责"读",写入目标(数据库事务、rollback 语义)由 destination provider 承担,编写 source 时无需关心回滚逻辑;
- 本文所有阶段名、方法名与调用时序均以当前仓库
packages/core/data-transfer源码为准,若升级到其他 Strapi 版本,建议重新核对 src/types/providers.ts 与 src/engine/index.ts 中的定义。
延伸阅读:Destination Providers 文档、Stream Lifecycle 文档、Engine 文档 分别覆盖了写流端契约、阶段流三阶段生命周期(创建/管道/关闭)与引擎选项全貌,与本文构成 data-transfer provider 文档的完整三件套。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0623
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00