03-全链路企业级算力防线蓝图 » 历史记录 » 修订 13
修订 12 (Huarui Lin, 2026-04-24 20:31) → 修订 13/14 (Huarui Lin, 2026-04-24 21:27)
# # 全链路企业级算力防线蓝图 (严格 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()` 后,提取每个 Segment 的绝对首行。若 `cumulative_net_value` 为 NULL,立即抛出 `core.exceptions.SegmentFirstRowNullError` 阻断流水线。
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 微线程 24 核微线程 -> 物理哨兵封堵”极限拓扑。
### Step 2.0: 2.1: S3 特征工程统一计算与落盘闭环 入口按年流式与局部安全池桥接
1. **宏观黑盒封装**:建立唯一的对外入口 `src/pipelines/s3/step_s3_feature_engineering.py`。统筹器仅需传入输入目录与输出目录,内部封装完整的“按年循环 -> 打散 -> 计算 -> 合并 -> 落盘”生命周期,严禁在统筹器中传递巨大的中间态内存字典。 单年数据全量拉入安全沙箱(< 3GB)。
2. **按年流式安全沙箱**:外层循环严格按年份读取 S2 产出的单年份物理文件(`s2_YYYY.parquet`),单次内存驻留强制锁定在 < 3GB 绝对安全红线内。
3. **局部豁免桥接打散**:在单年份安全池内,触发 执行 `partition_by(["fund_id", "segment_id"], as_dict=True)` 局部豁免权,将数据零拷贝打散为字典视图,作为向 as_dict=True)`,零拷贝打散为字典视图传递给 C 层传递连续内存的桥接器。 层。
### Step 2.2: 纯 NumPy C 算子微线程拓扑
1. 实例化 24 核 `ThreadPoolExecutor` 分发字典值。
4. **纯 NumPy 微线程分发**:实例化 `ThreadPoolExecutor`(核数依环境而定)分发字典值。GIL 2. GIL 防线:内层函数第一步必须 `to_numpy().flatten()` 抽出纯连续内存数组,后续计算强制约束为仅调用 `np.dot`、`np.sum` 等基础 抽出纯连续内存数组。
### Step 2.3: A/B/C 类特征与分母隔离防线
1. **分母隔离**:所有 12/26/52 周窗口计算,入口判断 `len(arr) < window`,不满足直接输出 `None`。
2. **A/B/C 类向量化**:使用 `cumsum` 截取法算均值偏离;`np.minimum(ret, 0)` 屏蔽正收益算下行波动;`cummax` 位移算最大回撤;使用 `is_start` 与 `cum_arr` 的 O(N) 骨架算连阴周数。
### Step 2.4: OLS 矩阵击穿与 EWM 封堵
1. **OLS 骨架**:提取 `np.log(nav)`,O(1) 公式算分母 `ss_xx`,纯 `np.dot` 算分子,斜率强制乘 `52` 年化。封杀 `lstsq`。
2. **EWM 骨架**:转换为 `scipy.signal.lfilter` 分子分母系数,利用 C 算子,利用底层短暂释放 GIL 实现真并行。 递推短暂释放 GIL。
5. **纯算子下沉与分母隔离**:所有 A/B/C/D/E 类特征计算逻辑作为纯函数下沉至 `src/core/math/` 供单测直接锤炼。所有带窗口期的特征计算,必须绝对遵从“分母隔离强制防线”:遇 3. **红线 3 拦截落地**:使用极小确定性输入(如 5 个固定净值点)运行 OLS/EWM 算子,生成精确至小数点后 10 位的 `.npy` 文件存入 `tests/fixtures/`。单测时强制比对当前输出与基线,误差越过 `1e-7` 直接阻断。
### Step 2.5: D/E 类特征与 NULL 必须直接输出 NULL,严禁残缺窗口计算。 穿透
1. **D 类定投**:分母强制修正为带窗口期的 `rolling_mean(nav, X)`,计算成本偏离与胜率。
6. **物理哨兵生成**:在 38 维特征合并后,立即执行全横向非空判定生成 `is_feature_complete` 列。 2. **E 类风险**:显式写出 `(mean_ret - weekly_rf) / std_ret` 结构。
7. **红线门禁与落盘**:复用全局 `PhysicalStorageContract.sink_parquet_strict()` 3. **Regime NULL**:上市不足 52 周导致分母为 NULL 时,直接输出 NULL,不做任何 `fill_null(0)` 替换,任由其触发哨兵。
### Step 2.6: 物理哨兵生成与落盘
1. 横向非空判定生成 `is_feature_complete`。
2. 复用全局 `sink_parquet` 拓扑产出 `s3_YYYY.parquet`。落盘后立即调用 `validate_schema_symmetry()` 比对文件 `s3_YYYY.parquet`。
3. **红线 2 拦截落地**:落盘后立即比对文件 Schema 与 S3 硬编码契约(3 列物理主键 硬编码契约(S2 基础列 + 38 维特征列 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) `is_feature_complete`)的对称差,差集非空直接阻断循环并报错。
---
## 第 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 是纯矩阵运算无中间态膨胀,其内存曲线必须是一条绝对水平的直线。任何呈阶梯状上升的现象,都证明存在隐藏的内存泄漏或未触发全局缓存的灾难逻辑,必须阻断上线。