Integration 架构
当前 Integration 数据链路分为三种形态:API-based、file-based ELT 和仍在迁移中的 legacy file-based。这里的 file-based 描述的是供应商交付方式;它不再意味着 Retail API 一定直接读取并解析文件。
flowchart LR
API[Vendor API] --> RA[Retail API] --> DB[(PostgreSQL)]
SFTP[Vendor SFTP / RC SFTP] --> C[Collector]
C --> RAW[(Raw S3)] --> L[File-based Integration Lambda]
L --> N[(Normalized S3)] --> Q[SQS]
Q --> E[integrations/elt]
E --> SF[(Snowflake target tables)]
E -->|completion SNS| RA
SF --> RA
LEGACY[Legacy file-based] -. direct file read .-> RA
File-based ELT 链路(已迁移 vendor)
Section titled “File-based ELT 链路(已迁移 vendor)”Collector从 vendor SFTP 拉取文件,或接收推送到 RC SFTP 的文件,并写入 raw S3。file-based-integration-lambda只负责 vendor-specific 的解压、解析和标准化,输出 normalized CSV 到 S3;不写业务数据库。- Lambda 发送
LoadRequest到 SQS。请求包含 vendor、as_of_date、dataset 路径和可选 tenant。 integrations/elt使用 SnowflakeCOPY INTO加载 staging,再执行 SQL transform,生成 vendor-agnostic target tables。- ELT 完成后通过 SNS 通知 Retail API。Retail API 根据 vendor 和 tenant 触发对应 integration 的 nightly sync,从 Snowflake target tables 读取结果并写入 PostgreSQL。
ELT 把 vendor 文件格式和业务数据库解耦:Collector/Lambda 处理输入格式,ELT 处理 schema mapping 与业务规则,Retail API 只消费稳定的 target contract。
Legacy file-based 链路
Section titled “Legacy file-based 链路”尚未完成 cutover 的 vendor 仍使用:
Vendor SFTP → Collector → S3 latest/文件 → Retail API parser → PostgreSQL这条链路通常依赖 latest/ 或 LATEST 文件,并由 Retail API 内的 file-based parser 完成读取、mapping 和写库。注意:部分已完成 Retail API ELT cutover 的 vendor,Lambda 仍可能为了兼容或其他下游生成 latest/LATEST;因此 LATEST 是否存在不能单独判断 Retail API 的消费链路,最终状态以 isMigratedToElt() 和当前 vendor 配置为准。
当前迁移状态
Section titled “当前迁移状态”截至 2026-08-20,代码中的 IntegrationType::isMigratedToElt() 明确标记以下 12 个 vendor 已切换到 ELT:
| 状态 | Vendor |
|---|---|
| 已完成 Retail API cutover | Altruist、Apex、Betterment、Flourish、Folio Investing、LPL、My529、Pacific Life、Raymond James、RBC、Schwab、Trust America |
| 仍需按 legacy 方式理解 | Fidelity、First Clearing、Interactive Brokers、Jackson、Pershing、SEI,以及其他尚未被 isMigratedToElt() 标记的 file-based vendor |
| 独立迁移路径 | Morningstar:Equity Data 使用 Snowflake Data Share;Fund/Price 文件进入 Collector + Lambda + ELT,详见 Morningstar 迁移文档 |
这张表以 retail-api 的 isMigratedToElt() 为准;integrations/elt/declarations/ 可能先于生产 cutover 提交,因此不能单独作为上线状态来源。
API-based 链路
Section titled “API-based 链路”常规 API-based integration 由 Retail API 直接调用 vendor API:
- Advisor 完成 OAuth 或其他授权流程。
- Retail API 保存加密 credentials 和 integration mapping。
- Sync job 调用 vendor API,转换响应并写入 PostgreSQL。
但“API-based”不是全局的执行方式保证;当前代码中也存在将 API-based nightly sync 的结果写入 S3、通过 SQS 交给 integrations/elt 并落到 Snowflake 的 vendor/invocation。排查时必须按具体 vendor contract、invocation 和运行配置确认链路,不能仅凭 integration 类型下结论。
SSO 架构
Section titled “SSO 架构”SSO 是独立于数据同步的 integration 类型:
flowchart LR
VP[Vendor Portal] --> SSO[SSO Controller]
SSO --> AUTH[Auth Service]
AUTH --> APP[RightCapital Application]
常见方向包括 Vendor → RightCapital、RightCapital → Vendor,以及带 household context 的 contextual SSO。SSO 不经过 Collector、Lambda 或 ELT。
Nightly Sync 与监控
Section titled “Nightly Sync 与监控”Nightly Sync 通常在 US market close 后运行,但不同 vendor 的文件到达时间可能不同。监控应覆盖:
- ELT vendor:raw/normalized 文件到达、SQS 消费、Snowflake target 更新、completion SNS、Retail API sync;
- legacy vendor:Collector、
LATEST文件、parser、queue job 和数据库更新; - API-based vendor:credentials、token refresh、vendor API error/rate limit 和 sync job。
具体排查步骤见 Nightly Sync 和 Data Issues。