Skip to content

Collector

Collector 是独立的文件采集服务,负责从 vendor SFTP 下载文件,或处理推送到 RC SFTP 的文件,并把 raw 文件写入 S3。它是文件型 integration 的入口,但不再是所有文件型 integration 的最终处理者。

Repository: gitlab.rightcapital.io/integrations/collector

flowchart LR
    V[Vendor SFTP / RC SFTP] --> D[Collector download/receive]
    D --> R[(raw S3)]
    R --> L[file-based-integration-lambda]
    L --> N[(normalized S3)]
    N --> E[integrations/elt]

Collector 的输出是 raw S3 object。它会在本地临时目录完成下载和基础处理,并用日期/timestamp 过滤已处理文件;成功上传后清理临时文件。对已迁移到 ELT 的 vendor,后续由 Lambda 标准化并触发 ELT;Collector 不创建 LoadRequest,也不写 Snowflake 或 PostgreSQL。

阶段 作用 说明
Download/Receive 从 vendor SFTP 获取文件,或接收 RC SFTP push 使用 vendor-specific credentials 和 file pattern,并按日期/timestamp 只取新文件
Process 解压 ZIP/GZ、校验格式、按类型整理 删除已成功展开的压缩包;不做业务 mapping
Upload 写入 raw S3 按结构化 S3 key 保存,并更新 timestamp.txt 等处理记录
Finalize 为仍需要该兼容输出的 vendor 更新 LATEST/latest 是否由 Retail API 消费必须按 isMigratedToElt() 和 vendor 配置确认
raw S3 → file-based-integration-lambda → normalized S3 → SQS → integrations/elt

Lambda 负责 vendor-specific normalize;ELT 负责 Snowflake staging、SQL transform 和 target tables;Retail API 在 completion SNS 后消费 target tables。

raw S3 → LATEST/latest → Retail API parser → PostgreSQL

尚未完成 Retail API ELT cutover 的 vendor 由 Retail API 直接读取并解析文件。部分已 cutover vendor 仍可能生成 LATEST 供兼容或其他下游使用,因此排查时不能只凭 LATEST 存在与否判断最终消费者。

S3 key 是 vendor-specific,不能用一个包含 tenant 的模板覆盖所有 vendor。代码中可见的典型形态包括:

raw/self-hosted/fidelity/{basename}
raw/self-hosted/pershing/{basename}
raw/vendor-hosted/first_clearing/{vendor_path}/{basename}
raw/vendor-hosted/morningstar/{yyyymmdd}/{file}
normalized/{vendor}/{tenant-or-date}/{as_of_date}/{dataset}.csv.CD.gz

实际 key 以 Collector configuration 和 Lambda vendor props 为准。并非所有 vendor 都按 Rep Code 分目录;有些 vendor 使用整份文件,有些使用 advisor、firm、RIA 或其他 tenant key。

Collector 处理的常见文件包括:

文件类型 内容 常见格式
Accounts account number、type、status、owner CSV / fixed-width
Positions security、quantity、market value CSV / fixed-width
Securities CUSIP、ticker、name CSV / fixed-width
Tax lots cost basis、acquisition date CSV / fixed-width
Transactions trade history CSV / XML / fixed-width

Collector 只负责文件层面的下载、解压和存储;字段含义、mapping 和业务校验由 Lambda 或 legacy parser 负责。

每个 vendor 的配置通常包括:

  • SFTP host、port 和 credentials;
  • 文件名 pattern、目录和 delivery timing;
  • host type(vendor-hosted 或 self-hosted);
  • 解压、重命名和目标 S3 path 规则;
  • ELT vendor 的 normalized dataset 触发条件(由下游 Lambda/配置实现)。

新增或修改配置时,先确认 vendor 的 Retail API 迁移状态、Lambda 的 skip_normalized_trigger 等配置,再决定是否需要 LATEST/finalize。迁移状态和 finalize 状态是两个独立维度。

  • SFTP 连接和 credentials 是否有效。
  • vendor 文件是否按预期到达,raw S3 是否生成 object。
  • Lambda 是否成功生成 normalized 文件(ELT vendor)。
  • legacy vendor 的 LATEST 是否更新。
  • 文件到达告警应按 vendor 迁移状态选择后续检查链路。
  • 认证失败、文件损坏、格式变化和临时目录清理失败都应保留原始错误上下文,便于重跑和定位。
现象 ELT vendor 检查 Legacy vendor 检查
没有新数据 raw object、Lambda normalize、SQS/ELT raw object、LATEST、Retail API parser
格式错误 Lambda vendor normalize logs Retail API parser logs
数据过期 as_of_date、Snowflake target、completion SNS LATEST 时间戳、sync job
tenant 对不上 LoadRequest tenant 和 ELT target Rep Code/reference 和 S3 路径
认证失败 核对 credentials、host type 和最近的 credentials rotation 核对 SFTP credentials、host、port 和 key/password
文件损坏 检查原始 object、解压日志和文件大小;必要时重新触发 Lambda 要求 vendor 重新发送文件,保留损坏文件和 run 证据
格式变化 对照 vendor schema、header 和 delimiter,更新 vendor-specific normalize 对照 parser 支持的 schema,确认 vendor 是否变更格式

排查时建议保留以下证据:vendor、文件名、delivery date、S3 key、文件大小、Collector run、Lambda/ELT run key。ELT vendor 还应记录 normalized dataset 是否完整;legacy vendor 则记录 LATEST 指向的实际文件和更新时间。