[Draft] DEV-19911 Migrating Morningstar to the ELT Pipeline
- DEV Ticket:DEV-19911 Morningstar: Equity Data Migration
- Tracking Ticket:ENGR-11632 Morningstar: Equity Data Migration September 30, 2026
- Related:OPS-1652 Obtain Snowflake Account ID for sharing with Morningstar
- Integration ELT Refactor
- Morningstar Equity Data Upgrade Hub
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:风险
Part 1:背景与需求
Section titled “Part 1:背景与需求”一、为什么现在必须做迁移?
Section titled “一、为什么现在必须做迁移?”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 同时调整了数据结构,因此本项目既要完成新数据接入,也要确认下游兼容性。
二、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 | 数据 | 接入方式 | 说明 |
|---|---|---|---|
| FTP | Fund Data、Equity Price | collector 通过 rclone 拉到 S3 | 不在本次 Equity Data Upgrade scope 内,依旧按原文件方式交付;后续纳入 collector + lambda 链路 |
| Snowflake Data Share | Equity Data | Snowflake Private Listing | Equity Data 直接通过 Snowflake 查询 |
| Crypto API | Crypto Symbol + Price | retail-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 路径,主要有三个原因:
- 实时性和执行效率更好:Snowflake Data Share 让 Morningstar 数据以可查询 view 的形式直接出现在 Snowflake 里,不需要等待文件下载、解压、解析和落库,后续 transform 可以更快开始。
- 文件方式处理复杂度更高:如果选择 SFTP,我们仍然需要处理结构复杂且体积巨大的 XML 文件(数 GiB 级别),还要设计全量 / 增量数据的下载、解析、比对和落库流程。Snowflake Data Share 把”数据已经可查询”作为起点,能把工程重心放在 schema mapping、转换上。
- 符合统一 ELT 方向:Snowflake 是我们 ELT 的底座。Morningstar 支持 Snowflake Data Share,正好和我们把 vendor 数据处理收敛到统一 ELT pipeline 的方向不谋而合。
Part 2:技术方案
Section titled “Part 2:技术方案”一、整体架构
Section titled “一、整体架构”拆分 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 内,仍按原文件方式交付,因此放进这条链路。
需要做的事情:
- 新增 Morningstar cron job,配置 Morningstar SFTP host、user、port、password / private key。
- 用 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。
需要做的事情:
- 新增
morningstarvendor pipeline。 - 为 Morningstar raw file 定义 input / output path contract:从
raw到normalized。 - 输出 gzip CSV,字段名和文件粒度与 ELT staging table spec 对齐。
- 在 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。
需要做的事情:
- 在
declarations/vendors.yaml新增morningstarvendor。 - 为 file-based source 定义 staging table YAML。
- 为 Snowflake share source 定义 transform 输入。
- 定义 target tables。
- 用 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 实现。
需要做的事情:
- 使用现有
snowflake.integrationconnection 读取INTEGRATION_ELT_*Snowflake database。 - 为 Morningstar 新增 sync。
- 将 Snowflake target rows upsert 到:
securities、equities、funds、prices、allocations_by_*、fund_*_portfolios等。 - 处理 delete / disable / restore / cache invalidation 语义,替代旧 Morningstar app 的 webhook + cache flush 行为。
- 在 retail-api 中直接同步 crypto symbol / price。
Part 3:实施计划
Section titled “Part 3:实施计划”一、工作项总览
Section titled “一、工作项总览”| 阶段 | 主要任务 |
|---|---|
| 1. 字段 Mapping 对齐 | 和 Morningstar 确认 Snowflake Data Share 配置;对照新版文档和现有系统字段,逐项确认字段含义、格式、数据类型和命名变化,判断旧字段是否能在新 schema 中找到对应或替代来源,并整理缺失字段和兼容性风险。 |
| 2. FTP collection | Fund 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、关键缺口已知,再进入正式开发。
Phase 2:开发
Section titled “Phase 2:开发”开发阶段按数据形态拆分实现:
- 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 实现;
Phase 3:双跑验证
Section titled “Phase 3:双跑验证”上线前新旧 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 双跑 |
Part 4:风险
Section titled “Part 4:风险”Deadline 风险
Section titled “Deadline 风险”Morningstar 明确当前 Equity Data delivery Deadline 2026-09-30。建议至少预留三周双跑窗口,用旧 pipeline 和新 ELT pipeline 做 diff。