Skip to content

[Draft] DEV-19911 Migrating Morningstar to the ELT Pipeline

Morningstar Morningstar 是一家全球性的金融数据和投资研究公司,是我们的证券数据供应商。我们用它提供证券基础信息、价格、基金配置、持仓、风格、行业和类别等数据,支持证券搜索、持仓匹配和投资组合分析。

Snowflake Data Share Snowflake Data Share 是 Snowflake 原生的数据共享机制:供应商把指定的 database、schema、table 或 view 以只读方式授权给我们,我们在自己的 Snowflake account 里直接查询这些对象,但数据仍由供应商维护,不需要文件传输或复制。

ELT ELT 是 Extract、Load、Transform 的缩写,意思是先把原始数据取回来并加载到数据仓库,再在数据仓库里做转换。和先转换再入库的 ETL 相比,ELT 更强调保留原始数据、集中转换逻辑,方便追踪、重跑和排查问题。


💡 概括一下:Morningstar 要求我们在 2026-09-30 前完成 Equity Data 升级,旧交付方式届时下线。本项目会把现有独立的 Morningstar FTP pipeline 迁入统一 ELT 架构:Equity Data 通过 Snowflake share 接入;Fund / Equity Price 等仍走 FTP 的数据,则通过 collector + lambda + ELT 进入 Snowflake,再由 retail-api 同步回业务数据库。

  • Part 1:背景与需求
  • Part 2:技术方案
  • Part 3:实施计划
  • Part 4:风险

Morningstar 正在升级 Equity Data,旧交付方式会在 2026-09-30 后下线。我们必须在此之前完成迁移,否则当前依赖 Morningstar equity data 的相关链路会失去数据来源。

为什么 Morningstar 要做这次升级?

Morningstar 这次升级的核心不是单纯更换交付渠道,而是把过去几年在 equity data 上的投入产品化:他们扩展了数据采集范围,提升了数据颗粒度、标准化程度和准确性,也新增了更多经过清洗和归一化的数据集。与此同时,客户对数据接入的要求也在变化:有的客户希望把大批量数据加载到自己的 data warehouse,有的希望通过 API 支撑应用数据库或前端展示。因此 Morningstar 在 2025-2026 年推动客户迁移到新版 equity data,并同步提供更现代的交付方式,包括 data feeds、API 和 Snowflake Data Share。

更多信息参考:Morningstar Equity Data Upgrade Hub

这次迁移不是简单换一个下载地址。Morningstar 同时调整了数据结构,因此本项目既要完成新数据接入,也要确认下游兼容性。

二、Morningstar 数据支撑哪些业务场景?

Section titled “二、Morningstar 数据支撑哪些业务场景?”

Morningstar 数据是我们 security 数据的重要来源,主要用于:

  • 证券识别:把持仓中的证券标识匹配到正确证券;
  • 估值:用价格和数量计算持仓市值;
  • 投资分析:提供资产类别、行业、风格、基金类别、配置比例和投资组合持仓等分析维度。

三、当前 Morningstar pipeline 是怎么工作的?

Section titled “三、当前 Morningstar pipeline 是怎么工作的?”

当前 Morningstar 数据由独立项目 integrations/morningstar 处理。它是一个 Laravel 应用,拥有自己的 S3 bucket 和数据库,负责下载 FTP 文件、解析数据、写入自己的数据库,再同步到 retail-api DB。

当前主流程如下:

flowchart TB
    subgraph equity_fund["FTP Equity / Fund"]
        FTP["Morningstar FTP"] --> retrieve["file:retrieve → S3"]
        retrieve --> import["data:import"]
        import --> transform["data:transform"]
    end

    subgraph crypto_pipeline["Crypto API"]
        API["Morningstar Web Services API"] --> sync_symbols["crypto:sync-symbols"]
        sync_symbols --> fetch_prices["crypto:sync-prices"]
    end

    transform --> push["sync:push"]
    fetch_prices --> push
    push --> api_db["Retail API DB"]

这套设计已经稳定运行十多年,适合 Morningstar 过去的文件型交付方式,也真正做到了几乎免维护。这次迁移不是因为旧系统模式失效,而是 Integration Team 正在把 vendor 数据处理收敛到统一的 ELT pipeline 架构,后续可以更好地统一维护、观测和演进。

四、为什么纳入 Integration ELT 架构?

Section titled “四、为什么纳入 Integration ELT 架构?”

Integration Team 正在把 vendor 数据处理迁移到统一的 ELT pipeline:上游负责稳定收集和标准化原始数据,Snowflake 负责集中存储和声明式转换,retail-api 消费转换后的结果并写回业务数据库。

Morningstar 这次迁移正好落在这个架构转型窗口内。它既有供应商强制升级带来的时间压力,也提供了把独立 Morningstar pipeline 收敛到统一 vendor 数据处理链路的机会。

五、新的数据来源会变成什么?

Section titled “五、新的数据来源会变成什么?”

迁移后,Morningstar 数据源会从原来的 2 条拆成 3 条:

Source数据接入方式说明
FTPFund Data、Equity Pricecollector 通过 rclone 拉到 S3不在本次 Equity Data Upgrade scope 内,依旧按原文件方式交付;后续纳入 collector + lambda 链路
Snowflake Data ShareEquity DataSnowflake Private ListingEquity Data 直接通过 Snowflake 查询
Crypto APICrypto Symbol + Priceretail-api 直接调用 Morningstar Web Services API数据量较小,且是 API 型交互,不进入 collector / lambda / ELT 链路

六、关键取舍:Equity Data 选择 Snowflake Data Share 路径

Section titled “六、关键取舍:Equity Data 选择 Snowflake Data Share 路径”

Morningstar 也提供了新的 SFTP 交付路径供我们选择。我们选择 Snowflake Data Share 作为 Equity Data 路径,主要有三个原因:

  1. 实时性和执行效率更好:Snowflake Data Share 让 Morningstar 数据以可查询 view 的形式直接出现在 Snowflake 里,不需要等待文件下载、解压、解析和落库,后续 transform 可以更快开始。
  2. 文件方式处理复杂度更高:如果选择 SFTP,我们仍然需要处理结构复杂且体积巨大的 XML 文件(数 GiB 级别),还要设计全量 / 增量数据的下载、解析、比对和落库流程。Snowflake Data Share 把”数据已经可查询”作为起点,能把工程重心放在 schema mapping、转换上。
  3. 符合统一 ELT 方向:Snowflake 是我们 ELT 的底座。Morningstar 支持 Snowflake Data Share,正好和我们把 vendor 数据处理收敛到统一 ELT pipeline 的方向不谋而合。

拆分 Morningstar 数据链路,按数据形态接入:文件型数据和 Snowflake Data Share 进入 ELT;Crypto API 数据量较小,直接放在 retail-api 同步。

flowchart LR
    msFtp[Morningstar FTP<br/>fund / equity price]
    msShare[Morningstar Snowflake Data Share<br/>Equity Data]
    cryptoApi[Morningstar Web Services API<br/>crypto symbol / price]

    collector[collector<br/>rclone FTP files to S3]
    lambda[lambda<br/>normalize S3 files<br/>and write back to S3]
    elt[integrations/elt]

    targetSf[Snowflake<br/>target tables]
    retailApi[retail-api<br/>sync to DB / crypto API sync]
    retailDb[Retail DB]

    msFtp --> collector --> lambda --> elt --> targetSf
    msShare --> elt
    targetSf --> retailApi --> retailDb
    cryptoApi --> retailApi

分层后的职责是:

项目负责什么不负责什么
collector从 Morningstar FTP 下载 Fund / Equity Price 原始文件到 S3不解析业务内容,不做 transform
file-based-integration-lambda把原始文件标准化成后续 SQL 处理友好的 CSV 格式不写业务数据库,不承担最终业务 mapping
integrations/elt把 normalized S3 files load 到 Snowflake staging;读取 Snowflake Data Share;执行 SQL transform,产出 target tables不负责 retail-api 业务写库;不处理 Crypto API 同步
retail-api从 Snowflake target tables 读取转换结果并写入 DB;直接调用 Morningstar Web Services API 同步 crypto symbol / price不直接解析 Morningstar 原始文件,不承担 ELT transform

二、collector:把文件型数据拉回 raw bucket

Section titled “二、collector:把文件型数据拉回 raw bucket”

collector 当前已经有成熟的 rclone SFTP → S3, cron job 模式。Fund Data 和 Equity Price 不在 Morningstar 本次 Equity Data Upgrade scope 内,仍按原文件方式交付,因此放进这条链路。

需要做的事情:

  1. 新增 Morningstar cron job,配置 Morningstar SFTP host、user、port、password / private key。
  2. 用 rclone 把指定目录下的 Fund / Equity Price 文件复制到 integration bucket:raw/vendor-hosted/morningstar/...

三、file-based-integration-lambda:把 raw files 变成 SQL 友好的 normalized data

Section titled “三、file-based-integration-lambda:把 raw files 变成 SQL 友好的 normalized data”

Morningstar 的 Fund Data 包含 XML,price 文件也有不同压缩和结构。lambda 完成转换,输出稳定、可声明 schema 的 normalized files。

需要做的事情:

  1. 新增 morningstar vendor pipeline。
  2. 为 Morningstar raw file 定义 input / output path contract:从 rawnormalized
  3. 输出 gzip CSV,字段名和文件粒度与 ELT staging table spec 对齐。
  4. 在 normalize 完成后触发 ELT load request。

四、integrations/elt:把 Morningstar 数据沉淀为 Snowflake target tables

Section titled “四、integrations/elt:把 Morningstar 数据沉淀为 Snowflake target tables”

对 SFTP source 的文件,ELT 负责把 lambda 产出的 normalized CSV load 到 Snowflake staging tables,再执行 SQL transform。对 Snowflake Data Share 路径,ELT 需要读取 Morningstar shared views,把数据转换成我们自己的 target schema。

需要做的事情:

  1. declarations/vendors.yaml 新增 morningstar vendor。
  2. 为 file-based source 定义 staging table YAML。
  3. 为 Snowflake share source 定义 transform 输入。
  4. 定义 target tables。
  5. 用 SQL transform 处理 identity mapping、status、price、allocation、style / sector 等字段。

五、retail-api:同步 ELT 结果,并承接 Crypto API

Section titled “五、retail-api:同步 ELT 结果,并承接 Crypto API”

retail-api 读取 Snowflake target tables,并将转换结果 upsert 到 DB。

Crypto symbol + price 数据量较小,当前在 integrations/morningstar 中实现,需要将其迁移到 retail-api 通过 cron job 实现。

需要做的事情:

  1. 使用现有 snowflake.integration connection 读取 INTEGRATION_ELT_* Snowflake database。
  2. 为 Morningstar 新增 sync。
  3. 将 Snowflake target rows upsert 到:securitiesequitiesfundspricesallocations_by_*fund_*_portfolios 等。
  4. 处理 delete / disable / restore / cache invalidation 语义,替代旧 Morningstar app 的 webhook + cache flush 行为。
  5. 在 retail-api 中直接同步 crypto symbol / price。

阶段主要任务
1. 字段 Mapping 对齐和 Morningstar 确认 Snowflake Data Share 配置;对照新版文档和现有系统字段,逐项确认字段含义、格式、数据类型和命名变化,判断旧字段是否能在新 schema 中找到对应或替代来源,并整理缺失字段和兼容性风险。
2. FTP collectionFund Data 和 Equity Price 仍按原文件方式交付;在 collector 中新增 Morningstar rclone cron job,把原始文件复制到 S3 raw bucket。
3. Normalize新增 Morningstar vendor,补 XML-to-CSV normalize,触发 ELT。
4. Snowflake ELT新增 Morningstar declarations、staging / target tables、transform SQL。
5. Retail sync / Crypto API从 Snowflake 读取 target rows,写 DB;直接同步 crypto symbol / price。

Phase 1:字段 Mapping 对齐,数据前期初步验证

Section titled “Phase 1:字段 Mapping 对齐,数据前期初步验证”

重点确认:

  • Morningstar 新版文档中的字段,和现有系统字段是否能一一对应或找到替代来源;
  • 字段含义、格式、数据类型是否变化;
  • Snowflake Private Listing 可用,用真实数据实际对照确认;

这个阶段的目标是先做一轮数据前期验证,确认新数据源可访问、核心字段可 mapping、关键缺口已知,再进入正式开发。

开发阶段按数据形态拆分实现:

  • collector:新增 Morningstar rclone cron job,把 Fund Data / Equity Price 原始文件复制到 S3 raw bucket;
  • lambda:新增 Morningstar vendor,把原始文件转换成后续 SQL 友好的 normalized files,并触发 ELT;
  • elt:新增 Morningstar declarations、staging / target tables 和 transform SQL;
  • retail-api:从 Snowflake target tables 同步数据到 DB;迁移 crypto 实现;

上线前新旧 Morningstar pipeline 并行运行一段时间,用旧 pipeline 结果和新 ELT 结果做 diff。确认差异可解释、风险可接受后,再切换。

时间目标
当前 - 2026-07-12完成字段 Mapping 对齐和数据前期初步验证
2026-07-13 - 2026-08-09完成主要开发
2026-08-10提交 Code Review
2026-08-11 - 2026-08-31预留 Code Review、修改、合并和双跑前验收时间
2026-09-01开始新旧 Morningstar pipeline 双跑

Morningstar 明确当前 Equity Data delivery Deadline 2026-09-30。建议至少预留三周双跑窗口,用旧 pipeline 和新 ELT pipeline 做 diff。