Skip to content

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)”
  1. Collector 从 vendor SFTP 拉取文件,或接收推送到 RC SFTP 的文件,并写入 raw S3。
  2. file-based-integration-lambda 只负责 vendor-specific 的解压、解析和标准化,输出 normalized CSV 到 S3;不写业务数据库。
  3. Lambda 发送 LoadRequest 到 SQS。请求包含 vendor、as_of_date、dataset 路径和可选 tenant。
  4. integrations/elt 使用 Snowflake COPY INTO 加载 staging,再执行 SQL transform,生成 vendor-agnostic target tables。
  5. 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。

尚未完成 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 配置为准。

截至 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 integration 由 Retail API 直接调用 vendor API:

  1. Advisor 完成 OAuth 或其他授权流程。
  2. Retail API 保存加密 credentials 和 integration mapping。
  3. Sync job 调用 vendor API,转换响应并写入 PostgreSQL。

但“API-based”不是全局的执行方式保证;当前代码中也存在将 API-based nightly sync 的结果写入 S3、通过 SQS 交给 integrations/elt 并落到 Snowflake 的 vendor/invocation。排查时必须按具体 vendor contract、invocation 和运行配置确认链路,不能仅凭 integration 类型下结论。

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 通常在 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 SyncData Issues