Nightly Sync 夜间同步
Nightly Sync 负责把 vendor 的最新数据同步到 RightCapital。它不是所有 file-based vendor 的统一“读 S3 文件”任务:已迁移 vendor 由 ELT 完成文件处理和 Snowflake transform,Retail API 在 ELT 完成通知后才执行业务同步;API-based 和 legacy file-based 仍由各自的 Retail API sync 流程处理。
Nightly Sync 通常按 US market close 后的部署调度运行;具体时间、vendor 到达窗口和队列配置以当前环境的 Cron/Kubernetes 配置为准,不能用固定的“每天几小时”作为 SLA。
ELT SQS worker 默认 SQS_MAX_RECEIVE_COUNT=3。retryable failure 会保留消息等待 visibility-timeout retry;超过最大 receive count 后,worker 将该运行标记为 retry exhausted。terminal failure 不应无限重试。Retail API 的 batch sync 通过 vendor-specific command dispatch queue jobs,具体 queue 和并发配置以 command 实现为准。
按 integration 类型的链路
Section titled “按 integration 类型的链路”ELT file-based vendor
Section titled “ELT file-based vendor”Collector → raw S3 → Lambda normalize → SQS → integrations/elt→ Snowflake target → completion SNS → Retail API nightly sync → PostgreSQLELT worker 的一次运行包含:
- 校验
LoadRequest。 - 将 normalized datasets 通过
COPY INTO加载到 staging tables。 - 执行 SQL transform,生成或更新 target tables。
- 记录 audit 信息并按 retryable / terminal failure 分类。
- 通过 SNS 发布完成事件。
Retail API 收到 completion 后按 vendor 和 tenant 找到 active integrations,再 dispatch SyncIntegrationOnNightly。业务数据从 Snowflake target tables 读取并写入 PostgreSQL。
Legacy file-based vendor
Section titled “Legacy file-based vendor”Collector → S3 latest/文件 → Retail API parser → sync job → PostgreSQL该路径依赖 LATEST 或 latest/ 文件,Retail API 负责读取、解析、mapping 和写库。
API-based vendor
Section titled “API-based vendor”Retail API sync job → Vendor API → mapping → PostgreSQLELT 与 Nightly Sync 是相邻但不同的阶段:
| 阶段 | 主要职责 |
|---|---|
| ELT | normalized file input 的 staging、transform、Snowflake target materialization |
| Nightly Sync | 将已准备好的 vendor 数据按 integration mapping 同步到业务模型 |
| PostgreSQL | Retail API 的业务运行数据 |
ELT 运行按 vendor、dataset 和 as_of_date 隔离,支持安全重跑和 backfill。completion publishing 失败不应改变 ELT 主处理结果;retryable failure 留在队列中重试,terminal failure 记录后移除。
已迁移 vendor
Section titled “已迁移 vendor”按以下顺序检查:
- Vendor 文件是否进入 raw S3。
- Lambda 是否生成完整 normalized datasets。
- 是否发送并消费 SQS
LoadRequest。 - ELT staging/target 是否成功、
as_of_date是否正确。 - completion SNS 是否到达 Retail API。
- Retail API 是否 dispatch sync,PostgreSQL 是否更新。
Legacy vendor
Section titled “Legacy vendor”检查 Collector 下载、S3 LATEST 文件、parser 日志、integration mapping 和 sync job。
API-based vendor
Section titled “API-based vendor”检查 credentials、token refresh、vendor API 响应、rate limit 和 sync job retry。
通用运行指标
Section titled “通用运行指标”| 指标 | 关注点 |
|---|---|
| Completion rate | ELT completion、nightly sync completion 是否正常 |
| Failure rate | retryable 与 terminal failure 的比例 |
| Duration | ELT load/transform 和 Retail API sync 是否异常变慢 |
| Data freshness | target tables 或 PostgreSQL 的最后更新时间 |
| File arrival | vendor 是否按约定 delivery,文件大小是否异常 |
日志至少应能关联 vendor、as_of_date、tenant/reference、integration mapping 和 run key。ELT 还应关联 SQS message 与 Snowflake load/transform 结果。
| 问题 | 可能原因 | 首个检查点 |
|---|---|---|
| Sync 没有启动 | scheduler、queue worker 或 completion event 异常 | scheduler、SQS/SNS、Retail API webhook |
| 数据过期 | vendor 未交付、ELT 未完成或 legacy LATEST 未更新 |
按 vendor 类型检查对应链路 |
| 失败率升高 | vendor API、文件 schema、Snowflake transform 或 mapping 变化 | run log 和 vendor-specific error |
| 处理变慢 | 文件量、Snowflake load/transform、API latency 或数据库瓶颈 | 各阶段 duration |
手动恢复前先判断 vendor 类型:ELT vendor 优先重发同一 as_of_date 的 LoadRequest 或重跑 ELT,再触发 Retail API sync;legacy/API vendor 才使用其对应的 Retail API sync command。命令名称和参数以当前 retail-api console command 为准,不要把 legacy LATEST 恢复步骤用于 ELT vendor。
当前代码中可见的手动入口示例:
# ELT vendor:扫描 normalized 文件并模拟 S3 → SQS triggerphp artisan file-based-integration-lambda:trigger --vendor=lpl --since=yesterday --confirm
# 检查已迁移 vendor 的文件到达php artisan file-based-integration-lambda:file-arrival-check --vendor=lpl --as-of-date=2026-08-19
# legacy/API vendor:仅对确认存在且适用的 command 执行 batch syncphp artisan fidelity:sync# 具体 vendor command 以当前代码中的 `artisan list` 和 `command --help` 为准ELT trigger 只负责重新提交输入事件,不等于已经完成 Snowflake transform;恢复后仍需检查 ELT target、completion SNS 和 Retail API sync。
手动恢复原则
Section titled “手动恢复原则”不要对已迁移 vendor 直接执行“重新解析 LATEST 文件”的 legacy 操作。优先从同一 as_of_date 重新生成或重发 ELT LoadRequest,确认 Snowflake target 正常后再触发 Retail API sync。具体 vendor 的 tenant 规则以 EltTriggerSyncController 和 vendor declaration 为准。