项目

一般

简介

行为

# 全链路企业级算力防线蓝图 (严格 0-7 阶段版)

第 0 阶段:项目脚手架与 CI 基线

本阶段是全链路的“宪法”,不写任何业务逻辑,仅确立物理约束与拦截规则。

Step 0.1: 企业级目录拓扑与依赖隔离

  1. 空间隔离:建立 src/core(纯算子,零业务逻辑)、src/pipelines/s2_to_s5(离线训练流水线编排)、src/api/s6_inference(在线微服务)、src/backtest(回测对接层)、configs、tests 的强约束目录。src/core 内层严禁使用 from .xxx import *,4个核心模块必须建立独立子目录,且在各自的 __init__.py 中显式白名单导出,防止跨层逆向导入导致的 DRY 拓扑破坏。
  2. 依赖锁死:在 pyproject.toml 中精确锁定 polars, duckdb, lightgbm, numpy, scipy, fastapi, pydantic 的大版本号,彻底封杀上游破坏性更新引发的静默精度漂移。【依赖纯净化红线】 严禁引入实验追踪库(如 mlflow、weights&biases)与超参搜索库(如 optuna、ray[tune])。规约 S5.2 已在拓扑层面彻底否决超参搜索,此类依赖仅会污染 Docker 镜像层与增加无用编译时间。

Step 0.2: 配置中心统一收敛

  1. 将 D-1 基金类型白名单、MAD 裁剪系数 1.4826 与宽容度倍数 5.0、Spearman 阈值 0.90 全部沉淀于 config.yaml。
  2. 建立 core/config_loader.py,启动时做 Pydantic 强校验,缺失配置直接 sys.exit(1) 阻断。

Step 0.3: CI/CD 红线拦截器部署 (强制门禁)

在 Git 流水线与本地 Pre-commit 中植入四大静态/动态拦截器:

  1. Pandas 零容忍:基于现有 AST 语法树扫描器 check_forbidden_imports.py,发现 import pandas 直接报错阻断。
  2. Schema 对称差拦截 (已解封内聚):通过 Pytest 注册 schema_drift 标记,在 S2/S3/S4 物理落盘后立刻触发强校验,比对产出文件的列名集合与硬编码契约的对称差,差集非空直接 sys.exit(1)。
  3. 底层算子金色单测 (已解封内聚):通过 Pytest 注册 golden_precision 标记,在 Step 2.4 OLS/EWM 算子实现时生成极小确定性 .npy 基线,后续每次跑单测强制比对,浮点误差越过 1e-7 直接阻断。
  4. 工程治理正则扫描:新增纯 Python 脚本 check_engineering_governance.py。全局拦截违规 API 与魔法数字硬编码。内置注释过滤逻辑,防止误杀文档说明。

第 1 阶段:S2 - 数据血缘与清洗拓扑

面对 2800 万行原始净值,走“DuckDB 粗筛压榨 -> Polars 流式微操 -> Out-of-core 物理落盘”降维路径。

Step 1.1: DuckDB 纯粗筛算子防火墙

  1. 底层转码:Polars 读取原始 CSV,强制指定类型(fund_id 为 Categorical,日期为 Date32),落盘为原生 Parquet 数据湖。严禁 DuckDB 直接对抗 CSV。
  2. D-1 全局物理下推:在 DuckDB 内基于配置白名单过滤基础信息,提取合法 fund_id 集合。
  3. 周频无状态聚合:利用唯一合法的 SQL 实体(基于 1970-01-05 基准的纯数学整数除法对齐周一),Inner Join 后按周取组内 max(date) 对应净值。粗筛产物导出为极小的 s2_raw_weekly.parquet,关停 DuckDB。
  4. 统筹编排器物理接入:本步骤严禁作为孤立脚本直接被执行,必须作为标准入口参数被 scripts/run_offline_pipeline.py 极简统筹器按拓扑顺序显式调用。统筹器仅负责全局路径解析与上游异常阻断,严禁在此处包含任何业务计算逻辑,保证各 Step 模块的绝对逻辑闭包与可独立测试性。
  5. 结构化性能遥测挂载:本步骤的主执行函数必须强制挂载 core.telemetry.profile_step("S1.1") 装饰器。执行完毕后,必须向 stderr 独立输出一条包含 duration_sec(绝对耗时)与 peak_rss_mb(内存波峰)的标准 JSON 字符串。严禁高频采样 CPU,彻底隔离正常业务数据流,为全链路性能基线提供零侵入数据源。

Step 1.2: Polars 外层 Fund 隔离与长缺失切断

  1. S2 入口 Schema 强断言:读取 Step 1.1 产物后,必须强制比对物理列名(fund_id, net_value_date, cumulative_net_value)与数据类型(Utf8, Date, Float64),不一致立即抛出 AssertionError 阻断流水线。
  2. 长缺失切断判定:按 fund_id 分组,计算相邻两行真实数据周日期差值。当差值(周)> MAX_GAP_WEEKS + 1(即 > 5)时,使用累加逻辑(cumsum)为每个连续物理块打上递增的 segment_id。单基金内 segment_id 从 0 起始,各基金独立计算。严禁丢弃任何已有数据行。
  3. 结构化性能遥测挂载:主执行函数必须强制挂载 core.telemetry.profile_step("S1.2") 装饰器,执行完毕后向 stderr 输出含 duration_sec 与 peak_rss_mb 的标准 JSON。
  4. 统筹编排器物理接入:本步骤由 scripts/run_offline_pipeline.py 统筹器显式调用,接收 Step 1.1 产物路径,返回内存中的 Polars DataFrame(4 列:fund_id, net_value_date, cumulative_net_value, segment_id),供下游 Step 1.3 消费。

Step 1.3: 内层 Segment 绝对闭包与短缺失 5 步强制顺序

  1. 微型连续日历生成:按 ["fund_id", "segment_id"] 提取起止日期,利用 pl.int_ranges 与 pl.duration 纯数学加法生成无缺失周序列,与原始表执行 Left Join,并打上 is_original 标志。
  2. 净值纯净性扫描与动态切断:在 Segment 闭包内执行 forward_fill() 后,利用向量化掩码(is_nan | is_null | <=0 的 cum_max 传播)精准定位首个脏数据点,将该点及之后的行物理切断丢弃,保障下游 S4 数学纯度。被切出的脏段首行必然非法,自然消亡,不再触发全链路 FATAL 阻断。
  3. 血缘切断重算收益率:基于纯净 ffill 序列,使用 shift(1) 向量化重算全量 weekly_return。严禁复用任何上游历史收益率数据。
  4. 占位行收益率置 NULL:执行 when(is_original).then(weekly_return).otherwise(None) 掩码,将 is_original=False 的行收益率强制置为 NULL,彻底切断因 ffill 导致的 0.0 幻影收益向下游 rolling_std 的传播。
  5. 结构化性能遥测挂载:主执行函数必须强制挂载 core.telemetry.profile_step("S1.3") 装饰器。
  6. 统筹编排器物理接入:本步骤由统筹器显式调用,接收 Step 1.2 的内存 DataFrame,返回 6 列内存 DataFrame 供下游 Step 1.4 消费。

Step 1.4: 按年切片落盘与 Schema 红线门禁

  1. 年份集合提取:从 Step 1.3 产出的全量 DataFrame 中提取不重复年份并排序。
  2. 逐年切片落盘:循环内按年份 filter,调用 PhysicalStorageContract.sink_parquet_strict() 落盘为 s2_YYYY.parquet。内部自动注入物理排序(["net_value_date", "fund_id"])与压缩参数(zstd / row_group_size=100000)。
  3. 循环内就地释放:每次落盘后 del yearly_df,降低进程 RSS 基线,为下游 S3 按年流式加载提供物理前提。
  4. Schema 红线门禁:每次落盘后调用 PhysicalStorageContract.validate_schema_symmetry(),读取 Parquet footer metadata(零数据行加载),与 S2_OUTPUT_COLUMNS 做对称差校验,差集非空直接抛出 SchemaDriftError 阻断循环。
  5. 结构化性能遥测挂载:主执行函数必须强制挂载 core.telemetry.profile_step("S1.4") 装饰器。
  6. 统筹编排器物理接入:由统筹器显式调用,接收 Step 1.3 的内存 DataFrame 与输出目录路径,返回已落盘文件路径列表。
  7. 设计备忘:Step 1.4 的核心价值不在于拦截 S2 阶段自身的内存波峰(S2 全量仅 ~270 万行 × 6 列,Polars 稀疏编码下约 30-50 MB),而在于将 S2 产出物理切分为按年分区文件,防止下游 S3 特征工程(38 维浮点列 × 270 万行)一次性全量加载导致 RSS 飙升至 5-8 GB 的真实雪崩。

第 2 阶段:S3 - 特征工程宏观黑盒

全链路算力雪崩最高危区。严禁将本阶段拆分为多个独立 Step 文件串联,必须走“宏观内聚黑盒 -> 按年安全池 -> 纯 NumPy 微线程 -> 物理哨兵封堵”极限拓扑。

Step 2.0: S3 特征工程统一计算与落盘闭环

  1. 宏观黑盒封装:建立唯一的对外入口 src/pipelines/s3/step_s3_feature_engineering.py。统筹器仅需传入输入目录与输出目录,内部封装完整的“按年循环 -> 打散 -> 计算 -> 合并 -> 落盘”生命周期,严禁在统筹器中传递巨大的中间态内存字典。
  2. 按年流式安全沙箱:外层循环严格按年份读取 S2 产出的单年份物理文件(s2_YYYY.parquet),单次内存驻留强制锁定在 < 3GB 绝对安全红线内。
  3. 局部豁免桥接打散:在单年份安全池内,触发 partition_by(["fund_id", "segment_id"], as_dict=True) 局部豁免权,将数据零拷贝打散为字典视图,作为向 C 层传递连续内存的桥接器。
  4. 纯 NumPy 微线程分发:实例化 ThreadPoolExecutor(核数依环境而定)分发字典值。GIL 防线:内层函数第一步必须 to_numpy().flatten() 抽出纯连续内存数组,后续计算强制约束为仅调用 np.dot、np.sum 等基础 C 算子,利用底层短暂释放 GIL 实现真并行。
  5. 纯算子下沉与分母隔离:所有 A/B/C/D/E 类特征计算逻辑作为纯函数下沉至 src/core/math/ 供单测直接锤炼。所有带窗口期的特征计算,必须绝对遵从“分母隔离强制防线”:遇 NULL 必须直接输出 NULL,严禁残缺窗口计算。
  6. 物理哨兵生成:在 38 维特征合并后,立即执行全横向非空判定生成 is_feature_complete 列。
  7. 红线门禁与落盘:复用全局 PhysicalStorageContract.sink_parquet_strict() 拓扑产出 s3_YYYY.parquet。落盘后立即调用 validate_schema_symmetry() 比对文件 Schema 与 S3 硬编码契约(3 列物理主键 + 38 维特征列 + 1 列哨兵)的对称差,差集非空直接阻断循环并报错。
  8. 内存就地释放:年份计算完毕落盘后,立即 del 释放年份内存池,进入下一年迭代。
  9. 结构化性能遥测挂载:主执行函数必须强制挂载 core.telemetry.profile_step("S3_Feature_Engineering") 装饰器。

附录:S3 产出 42 列 Schema 物理死锁契约

特征列名必须与 0-architecture_baseline.md 第 4.3 节严格一字不差,严禁擅自增删改。

  • 3 列物理主键:fund_id, net_value_date, segment_id
  • 38 维 Float64 特征列:
    • A类 均值回归 (5):price_vs_ma_ratio_12w, price_vs_ma_ratio_26w, price_vs_ma_ratio_52w, price_vs_ma_ratio_ewm_26w, price_vs_ma_ratio_ewm_52w
    • B类 波动 (7):weekly_return, rolling_std_12w, rolling_std_26w, rolling_std_52w, downside_vol_26w, downside_vol_52w, volatility_regime
    • C类 趋势动量 (11):momentum_4w, momentum_12w, momentum_26w, momentum_52w, max_drawdown_26w, max_drawdown_52w, trend_slope_26w, trend_slope_52w, trend_r_squared_26w, trend_r_squared_52w, consecutive_down_weeks
    • D类 定投 (10):dca_cost_ratio_12w, dca_cost_ratio_26w, dca_cost_ratio_52w, dca_return_12w, dca_return_26w, dca_return_52w, dca_return_vol_26w, dca_return_vol_52w, dca_win_rate_26w, dca_win_rate_52w
    • E类 风险调整 (5):rolling_sharpe_26w, rolling_sharpe_52w, rolling_sortino_26w, rolling_sortino_52w, calmar_ratio_52w
  • 1 列物理哨兵:is_feature_complete (Boolean)

第 3 阶段:S4 - 标签生成

业务逻辑核心,防“毒特征污染”与“索引错位”。

Step 3.1: S4 绝对隔离加载

  1. scan_parquet 扫描 S2 产物,立即 .select(["fund_id", "net_value_date", "segment_id", "cumulative_net_value"])。严禁拉入 weekly_return。

Step 3.2: DCA 闭式向量化公式

  1. 断言 nv_array 纯净且大于 0。
  2. 倒数累加范式:inv_nv -> cum_inv -> weeks -> returns,全程无循环。

Step 3.3: 止盈边界与 0-based 转 1-based 拦截

  1. np.where(returns >= 0.20)[0] 命中检索。
  2. 索引偏移修正:label = int(hit_indices[0] + 1) if len > 0 else 0。
  3. 绝对边界断言:assert 0 <= label <= 150。

Step 3.4: 按年落盘

  1. 复用全局工具函数产出 s4_YYYY.parquet。
  2. 红线 2 拦截落地:落盘后立即比对文件 Schema 与 S4 硬编码契约 {"fund_id", "net_value_date", "segment_id", "label"} 的对称差,差集非空直接阻断循环并报错。

第 4 阶段:S5 - 标准化与 Spearman 降维

融合“截面数学推导”与“内存防爆累加”。

Step 4.1: 截面聚合与双形态参数物理时序隔离

  1. 纯横截面聚合:按 net_value_date 分组计算 38 维的 mean, std, median, mad,推导 MAD 上下界,生成 is_production_ready 哨兵。
  2. 路径 A (训练宽表):宽表上 .shift(1) 实现 T-1 对齐,按日期排序落盘。
  3. 路径 B (推理长表):Unpivot 为长表,动态生成自然日历 forward_fill 至每日,按 ["feature_name", "date"] 绝对物理排序落盘。

Step 4.2: 标准化 Apply 与 IEEE 754 拦截

  1. 实体表 .collect() 后与几百行宽表参数 Join(零内存膨胀广播)。
  2. 遍历 38 列构建 when(std==0).then(None).otherwise(...) 表达式,封杀除零产生 Inf。

Step 4.3: 训练 Join 三键死锁与毒样本截断

  1. 三键强锁 Join:["fund_id", "net_value_date", "segment_id"]。
  2. 哨兵覆写:is_feature_complete == False 的行,标签强制覆写为 0。

Step 4.4: Spearman 秩相关内存防爆防线

  1. 按年流式加载标准化特征。
  2. 截面内 rankdata(method="average") 锁定平均秩,算秩后立即减均值。
  3. Pearson 展开项 np.dot 联合累加 sum_xy,利用对称性减半算力。
  4. 跨年累加后执行 (matrix + matrix.T) / 2.0 强制对称化截断。
  5. 剔除 >= 0.90 的冗余特征,输出 final_feature_list。

第 5 阶段:S5.3 - 训练流水线

将降维后的纯净特征转化为 LightGBM LambdaRank 可消费的形态并进行闭环重训。

Step 5.1: LambdaRank 样本平衡重构

  1. Group 划分:全量表按 net_value_date 排序,每个日期即一个 Query Group。
  2. 差异化等距抽样:正样本池(label>0)用 np.linspace 抽最多 50 条;负样本池抽最多 15 条。
  3. 相关性得分映射:label=0 映射为 0,label=1~150 映射为 151 - label(第 1 周得 150 分)。
  4. Group 边界构建:记录每个截面抽样数,构建一维数组传给 LightGBM 的 group 参数。

Step 5.2: Time-Series 5 折评估防线

  1. 时间切分:按 20% 步长划分 5 个验证区块,使用累积扩展窗口。
  2. NDCG@10 绝对核心:训练传入 group_boundaries,评估强制锁定 NDCG@10 指标,直接对齐 Top 10 榜单业务 KPI。
  3. 防穿透锁定:5 折期间严禁再次剔除特征或修改任何超参,仅观察指标分布。

Step 5.3: 闭环全量重训与产物契约落盘

  1. 使用全量抽样数据 + final_feature_list 进行最终无截断训练。
  2. 产物强制打包:将 model.txt 与两张标准化参数表打入同一个带时间戳的发布目录(如 artifacts/v_20231024/)。
  3. MD5 强校验:生成模型文件的 MD5 摘要,作为后续 S6 加载前的防篡改门禁。

第 6 阶段:S6 - 推理服务

纯粹的“无状态数学计算引擎”,彻底剥离所有 I/O 职责。

Step 6.1: FastAPI 骨架与进程级全局缓存防抖

  1. 生命周期加载:在 lifespan 启动钩子中,读取发布目录产物,加载 _PARAMS_LONG_CACHE (长表)、_AVAIL_DATES_CACHE (日期数组)、_MODEL (Booster)。
  2. MD5 门禁校验:加载模型前比对 MD5,不一致直接 os._exit(1) 拒绝启动。

Step 6.2: 两步走时间锚定防线

  1. 第一步锚定:接收单维 request_date,与 _AVAIL_DATES_CACHE 执行极速 asof_join 寻找 valid_date。无效则短路返回空结构体。
  2. 第二步拉取:基于 final_feature_list 构建左表,按 ["feature_name", "valid_date"] 排序,与 _PARAMS_LONG_CACHE 执行双键 asof_join,后置过滤 is_production_ready。

Step 6.3: 2D 矩阵级向量化批处理内核

  1. 统一抽象:无论单基金(N=1)还是批量(N=1000),统一接收 (N, 38) 原始特征 2D 矩阵。
  2. 复用 S5 内核:调用 core.standardization_math.apply_zscore,利用 NumPy 广播一次性完成 N 行的 MAD 裁剪与 Z-Score 转换。
  3. 矩阵预测:_MODEL.predict(matrix),利用 C++ 底层极速输出得分。

Step 6.4: 接口物理隔离与归因红线

  1. 单基金接口:允许调用 predict_contrib 返回 SHAP 归因明细。
  2. 批量接口:架构级禁止引入 predict_contrib 代码,归因字段强制硬编码为空字符串 ""。

第 7 阶段:端到端集成与打通

这是将离线数学模型转化为线上业务价值的最终闭环,解决“产物交接、预热、回测对接、极限压测”四大工程鸿沟。

Step 7.1: 产物自动化交接与灰度发布契约

  1. CI/CD 打包流:S5 训练流水线成功后,CI 自动将 artifacts/v_xxx/ 打包为 Docker Image 或推送到 OSS/Minio。
  2. S6 滚动更新:K8s/ Docker Compose 执行滚动更新。新 Pod 启动时拉取最新产物进行 MD5 校验与预热,校验失败则 Pod 启动失败,旧 Pod 继续服役,实现无损回滚。

Step 7.2: S6 冷启动掩盖与 Health Check 强断言

  1. 就绪探针:在 FastAPI 中实现 /health 接口。不仅返回 200,内部必须验证 _MODEL is not None 且 _PARAMS_LONG_CACHE.shape[0] > 0。未加载完毕前,K8s 绝不将流量打入该 Pod。
  2. 预热请求:Pod 启动后,后台线程自动发起一次包含 10 只基金的矩阵预测请求,强制触发 LightGBM 内部树的 JIT 缓存预热,掩盖首次真实请求的微秒级延迟毛刺。

Step 7.3: 回测引擎标准化数据对接

  1. 定义回测接口契约:S6 批量接口不直接服务 C 端,而是作为 BFF 层被回测引擎(如 Backtrader/Qlib)调用。
  2. 历史截面回放:回测引擎按历史日期循环,将历史 T 日全市场基金原始特征组装为 (N, 38) 矩阵,调用 S6 获取 T 日得分排名。
  3. 严格按照 NDCG@10 验证:回测引擎拿到 Top 10 池后,向后看 150 周计算真实定投收益率。对比回测出的“Top 10 平均达标率”与 S5 验证集的“NDCG@10 曲线”,两者必须呈现强正相关,否则证明存在数据穿越 Bug。

Step 7.4: 全链路极限压测与 cgroup 防雪崩验证

  1. 流量炮台构建:使用 Locust 编写压测脚本,模拟 50 并发线程,每线程持续发送 N=1000 的批量推荐请求。
  2. P99 延迟红线:监控 S6 服务的 P99 响应时间,必须死锁在 < 50ms 以内。一旦超限,说明触犯了矩阵广播外的隐性锁或内存交换。
  3. 内存防雪崩观察:在压测期间,通过 docker stats 严格监视 S6 容器的内存曲线。因为 S6 是纯矩阵运算无中间态膨胀,其内存曲线必须是一条绝对水平的直线。任何呈阶梯状上升的现象,都证明存在隐藏的内存泄漏或未触发全局缓存的灾难逻辑,必须阻断上线。

由 Huarui Lin 更新于 5 个月 之前 · 14 修订