首页
/ Strapi Data Transfer 深入解析:Source Provider(ISourceProvider)接口与阶段流实现

Strapi Data Transfer 深入解析:Source Provider(ISourceProvider)接口与阶段流实现

2026-09-04 15:16:29作者:虞亚竹Luna

本文以 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 的方法分三类:

分类 方法 是否必需 作用
身份标识 typename 必需(类型字段) type 固定为 'source'name 是 provider 的唯一名称
生命周期 bootstrap(diagnostics?)close()beforeTransfer() 可选 引擎在初始化、收尾、阶段执行前分别调用
元信息 getMetadata()getSchemas()getStageTotals(stage) getMetadata 必需,其余可选 提供版本信息、schema 与阶段总量预估
阶段读流 createEntitiesReadStream()createLinksReadStream()createAssetsReadStream()createConfigurationReadStream()createSchemasReadStream() 均为可选方法 每个阶段一个 Readable 流,逐个 write(entity) 输出数据

需要注意两点源码级事实:

  1. 引擎在构造时会对两个 provider 做结构校验,见 engine/index.ts#L152-L163 中的 validateProvider('source', sourceProvider),因此实现必须满足接口契约(尤其是 type === 'source')。
  2. 五个 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.tsTransferStageTypeMaptypes/common-entities.ts

阶段 stage 读流方法 流内每条数据的类型 说明
schemas createSchemasReadStream() Struct.Schema 内容类型/组件的 schema 定义
entities createEntitiesReadStream() IEntity 不含关联的内容数据:{ type, id, data }
links createLinksReadStream() ILinkIBasicLink | IMorphLink | ICircularLink 实体间关系,kindrelation.basic / relation.morph / relation.circular
assets createAssetsReadStream() IAsset 媒体文件:{ filename, filepath, stream, stats, metadata },其中 stream 是二进制 Readable
configuration createConfigurationReadStream() IConfigurationcore-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_STAGESengine/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/promisespipeline 串联 source → transform → tracker → destinationpipeline 会在上游流自然结束(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 的完整调用链梳理为:

  1. 构造期validateProvider('source', ...) 结构校验(L155);
  2. bootstrap():以 Promise.allSettled 并发执行 source 与 destination 的 bootstrap(this.diagnostics),source 侧可在此建立数据库连接、打开 Strapi 文件包、校验授权;任一侧 reject 会触发 panic(L705-L716);
  3. 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 若不提供版本元数据,版本一致性检查会被跳过
  4. integrityCheck():两侧 getSchemas() 提供的 schema 集合经 #assertSchemasMatching 做策略化 diff(schemaStrategy 默认 strict),source 的 schema 输出质量直接影响迁移能否通过校验;
  5. beforeTransfer():先 source 后 destination 依次执行,异常可被已注册的 error handler 接管;
  6. 每个阶段:先调 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)";
  7. 阶段执行#transferStage 中依次发出 stage::start → 每写入一条数据 stage::progressstage::finish(或 stage::error / stage::skip)事件,事件挂在引擎的 progress.stream(PassThrough)上;整体迁移另有 transfer::init/start/finish/error 事件;
  8. 收尾close() 并发调用双方 close?.();若迁移中途抛错,引擎会调用 destination 的 rollback 并抛出原错误,source 侧的读流由 pipeline 的错误路径负责销毁。

另外,shouldSkipStage(L563-L588)基于 TransferGroupPresetscontent = 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.streamengine/index.ts#L338-L393),因此 source 侧应保证 asset.stream 可被单一消费者 pipe,且 stats.sizemetadata 尽量完整(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.tslinks.test.tsassets.test.tsconfiguration.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.tssrc/engine/index.ts 中的定义。

延伸阅读:Destination Providers 文档Stream Lifecycle 文档Engine 文档 分别覆盖了写流端契约、阶段流三阶段生命周期(创建/管道/关闭)与引擎选项全貌,与本文构成 data-transfer provider 文档的完整三件套。

登录后查看全文
热门项目推荐
相关项目推荐