Skip to content

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 实现为准。

Collector → raw S3 → Lambda normalize → SQS → integrations/elt
→ Snowflake target → completion SNS → Retail API nightly sync → PostgreSQL

ELT worker 的一次运行包含:

  1. 校验 LoadRequest
  2. 将 normalized datasets 通过 COPY INTO 加载到 staging tables。
  3. 执行 SQL transform,生成或更新 target tables。
  4. 记录 audit 信息并按 retryable / terminal failure 分类。
  5. 通过 SNS 发布完成事件。

Retail API 收到 completion 后按 vendor 和 tenant 找到 active integrations,再 dispatch SyncIntegrationOnNightly。业务数据从 Snowflake target tables 读取并写入 PostgreSQL。

Collector → S3 latest/文件 → Retail API parser → sync job → PostgreSQL

该路径依赖 LATESTlatest/ 文件,Retail API 负责读取、解析、mapping 和写库。

Retail API sync job → Vendor API → mapping → PostgreSQL

ELT 与 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 记录后移除。

按以下顺序检查:

  1. Vendor 文件是否进入 raw S3。
  2. Lambda 是否生成完整 normalized datasets。
  3. 是否发送并消费 SQS LoadRequest
  4. ELT staging/target 是否成功、as_of_date 是否正确。
  5. completion SNS 是否到达 Retail API。
  6. Retail API 是否 dispatch sync,PostgreSQL 是否更新。

检查 Collector 下载、S3 LATEST 文件、parser 日志、integration mapping 和 sync job。

检查 credentials、token refresh、vendor API 响应、rate limit 和 sync job retry。

指标 关注点
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_dateLoadRequest 或重跑 ELT,再触发 Retail API sync;legacy/API vendor 才使用其对应的 Retail API sync command。命令名称和参数以当前 retail-api console command 为准,不要把 legacy LATEST 恢复步骤用于 ELT vendor。

当前代码中可见的手动入口示例:

Terminal window
# ELT vendor:扫描 normalized 文件并模拟 S3 → SQS trigger
php 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 sync
php artisan fidelity:sync
# 具体 vendor command 以当前代码中的 `artisan list` 和 `command --help` 为准

ELT trigger 只负责重新提交输入事件,不等于已经完成 Snowflake transform;恢复后仍需检查 ELT target、completion SNS 和 Retail API sync。

不要对已迁移 vendor 直接执行“重新解析 LATEST 文件”的 legacy 操作。优先从同一 as_of_date 重新生成或重发 ELT LoadRequest,确认 Snowflake target 正常后再触发 Retail API sync。具体 vendor 的 tenant 规则以 EltTriggerSyncController 和 vendor declaration 为准。