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 配置确认 |
两种下游模式
Section titled “两种下游模式”ELT vendor
Section titled “ELT vendor”raw S3 → file-based-integration-lambda → normalized S3 → SQS → integrations/eltLambda 负责 vendor-specific normalize;ELT 负责 Snowflake staging、SQL transform 和 target tables;Retail API 在 completion SNS 后消费 target tables。
Legacy vendor
Section titled “Legacy vendor”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 配置
Section titled “Vendor 配置”每个 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 是否变更格式 |
文件到达检查
Section titled “文件到达检查”排查时建议保留以下证据:vendor、文件名、delivery date、S3 key、文件大小、Collector run、Lambda/ELT run key。ELT vendor 还应记录 normalized dataset 是否完整;legacy vendor 则记录 LATEST 指向的实际文件和更新时间。