00-基金量化定投选基系统-设计规约v324 » 历史记录 » 版本 3
Huarui Lin, 2026-04-25 12:36
| 1 | 1 | Huarui Lin | # 基金量化定投选基系统 — 设计规约 v 3.2.4 |
|---|---|---|---|
| 2 | **document_id:** SPEC-FUND-DCA-001 |
||
| 3 | **version:** 3.2.4 |
||
| 4 | **status:** Frozen (Baseline) |
||
| 5 | **created_date:** 2026-04-11 (v1.0.0) | 2026-04-30 (v 3.1.0) | 2026-05-08 (v 3.2.0) | 2026-05-09 (v 3.2.1) | 2026-05-10 (v 3.2.2) | 2026-05-11 (v 3.2.4) |
||
| 6 | **approved_by:** Henry Lin |
||
| 7 | **redmine_link:** https://redmine.persys.top/projects/fund-selection-system |
||
| 8 | # 第一部分:业务算法与数据拓扑红线 |
||
| 9 | **受众**:量化研究员、算法工程师、产品经理。 |
||
| 10 | **核心原则**:语言无关。本篇仅定义数学公式、数据流转拓扑与 Schema 边界,严禁在此篇中规定任何特定编程语言的 API 调用(如 `partition_by`、`gc.collect()` 等)。 |
||
| 11 | ## 1. 业务决策冻结域 |
||
| 12 | | 编号 | 决策项 | 最终决策 | |
||
| 13 | | :--- | :--- | :--- | |
||
| 14 | | D-1 | 类型表无时间维度 | 使用最新 `fund_type` 全局过滤。被过滤基金的全量历史数据一并剔除,不做追溯。 | |
||
| 15 | | D-2 | 标签定义与止盈阈值 | DCA 简单平均成本法,**总绝对收益率**达到 20% 止盈,150 周窗口,未达则 0。 | |
||
| 16 | | D-3 | label=0 样本保留 | 单基金上限 15 条。 | |
||
| 17 | | D-4 | 无风险利率 | 当前版本全链路物理丢弃该参数,不参与计算,仅作预留接口。 | |
||
| 18 | | D-5 | 增量更新 | 不采用增量,每周全量刷盘重训。 | |
||
| 19 | | D-6 | 止盈贴近度特征 | 不添加。 | |
||
| 20 | | D-7 | IC 计算 | 严格限制在 Time-Series Split 训练折内局部计算。 | |
||
| 21 | * **【D-4 工程翻译覆写】**:虽 D-4 声明“物理丢弃不参与计算”,但在 E 类特征向量化广播计算时,**代码中必须保留 `weekly_rf` 变量实体并硬编码赋值为 `0.0` 参与减法运算**。严禁直接在公式中硬编码字面量 `0.0`,严禁省略减法步骤。 |
||
| 22 | * **备注原因**:规约第一部分为纯数学定义,必须保留符号以维持体系纯洁性。而在工程实施中,如果直接删除变量或写成字面量 `0.0`,未来一旦引入真实无风险利率,研发必须全局硬核搜索替换,极易漏改且破坏公式拓扑。保留变量赋 `0.0` 是锁定未来参数接入时“零代码结构变动成本”的唯一正确做法。 |
||
| 23 | ### 1.1 样本平衡机制细节(业务红线) |
||
| 24 | * **业务动作**:针对单基金的正负样本,必须执行差异化的上限截断。label > 0 单基金最多保留 50 条;label = 0 单基金最多保留 15 条。抽样策略必须采用等距抽样。 |
||
| 25 | * **备注原因**:定投场景下,“未达标(label=0)”的样本量天然远大于“达标”样本。若不进行差异化截断,模型会在训练时被海量负样本淹没,丧失对稀缺止盈信号的学习能力;等距抽样保证了被抽中样本在时间轴上的均匀分布,防止时间维度的过拟合。 |
||
| 26 | ## 2. 数据血缘与清洗拓扑 (S2) |
||
| 27 | ### 2.1 双层逻辑闭包拓扑 |
||
| 28 | 严禁在单基金全局范围内处理短缺失。必须强制采用**“外层按 Fund 隔离 + 内层按 Segment 隔离”**的双层绝对逻辑闭包。 |
||
| 29 | * 在最内层的单 Segment 闭包内,动态提取该 Segment 的起止日期,生成微型连续周频日历表,后续补齐仅在此微型表上执行。 |
||
| 30 | ### 2.2 数据粗筛与周频对齐铁律 |
||
| 31 | * **DuckDB 粗筛强断言**:一周内存在多条净值记录时,**必须取该周期内最大日期对应的 `cumulative_net_value` 作为该周的唯一代表值**。 |
||
| 32 | * **备注原因**:若采用均值或首日聚合,会严重扭曲该周的真实净值状态,导致后续所有基于净值衍生的收益率与动量特征全部失效。 |
||
| 33 | * **无状态 SQL 铁律**:严禁使用受系统时区环境变量影响的日期截断函数。必须强制使用纯数学偏移算法(基准日必须硬编码为绝对周一 `1970-01-05`)进行周频聚合。 |
||
| 34 | * **备注原因**:“基准日锁定为绝对周一”仅为业务话,若不给出具体的数学锚点,研发极易在不同数据库引擎间实现出看似无状态实则不一致的偏移函数。`1970-01-05`(Unix 纪元第 5 天)是保证跨引擎物理行为绝对一致的唯一确定性基准。 |
||
| 35 | * **长缺失切断的精确定义**:当检测到相邻两行真实数据的周日期差值绝对大于 `max_gap_weeks + 1` 时,严禁丢弃任何已有数据行。必须将当前连续块封闭,`segment_id` 递增。 |
||
| 36 | * **备注原因**:`+1` 的精确偏移是为了吸纳正常的跨周边界(如本周五到下周一算差 1),只有真正意义上的长期停复牌(如超过 5 周)才会触发物理切断,防止误杀有效数据段。 |
||
| 37 | * **【DuckDB 算力防火墙】**:粗筛阶段严禁执行任何滚动窗口计算,严禁使用序列生成函数补齐缺失日期。 |
||
| 38 | * **备注原因**:数据库引擎不具备向量化处理高频金融时序切片的能力,强行计算会导致算力雪崩与内存溢出,必须将纯粗筛结果拉入内存引擎处理。 |
||
| 39 | ### 2.3 【核心决策记录】短缺失占位行收益率强制置 NULL 的业务防御 |
||
| 40 | * **业务动作**:在短缺失占位行(Gap ≤ 4)完成 `cumulative_net_value` 的 Ffill 后,**必须将该占位行的 `weekly_return` 强制置为 `NULL`**。 |
||
| 41 | * **短缺失 5 步强制顺序**(在最内层单 Segment 闭包内执行):Gap ≤ `max_gap_weeks` 时: |
||
| 42 | 1. 插入占位行并标记内存中间态列 `is_original`; |
||
| 43 | 2. 仅对净值列执行 Ffill; |
||
| 44 | 3. **【首行 Ffill 越界防御】**:由于已处于绝对隔离的单 Segment 闭包内,必须直接触发极简断言:验证当前闭包首行的净值是否为非空。若首行仍为 NULL,**必须立即抛出异常阻断流水线**。 |
||
| 45 | * **备注原因**:如果某只基金在 Segment 起始就存在数据损坏(首行全空),不阻断流水线的话,后续不仅会向下 Ffill 出全量“幽灵数据”,还会被错误地当作合法数据输入模型,产生大规模毒特征污染。 |
||
| 46 | 4. 基于纯净净值序列重新计算全量收益率; |
||
| 47 | 5. 仅当“当前行本身”为占位行时,将该行收益率强制置 NULL; |
||
| 48 | 6. 物理删除中间态列后落盘。 |
||
| 49 | * **决策原因(严禁修改)**:如果保留 Ffill 导致的 `0.0` 收益率,下游 S3 特征工程在计算 `rolling_std(weekly_return, 12)` 等波动率特征时,`0.0` 会被视为合法有效样本参与方差计算。这将人为压低短缺失恢复期的波动率,产生“伪低波动”毒特征,严重误导模型对风险的感知。 |
||
| 50 | * **防御闭环**:置 NULL 后,完美触发 S3 特征工程的“分母隔离防线”(遇 NULL 直接输出 NULL),最终这些毒样本会被 S3 产出的物理哨兵列 `is_feature_complete` 捕获,并在 S5 训练流水线中被强制截断为 `label=0`。这是牺牲局部样本量以换取全局特征纯洁性的唯一正确拓扑。 |
||
| 51 | * **weekly_return 血缘切断防线**:严禁沿用上游粗筛阶段产出的收益率。必须在当前绝对隔离的 Segment 闭包内,基于纯净的净值序列重新计算。 |
||
| 52 | * **备注原因**:上游的收益率计算无法感知后续的 Segment 切断逻辑,沿用会导致切断点附近的收益率发生跨时空错乱。 |
||
| 53 | ### 2.4 物理存储与 Schema 契约 |
||
| 54 | * **S2 入口 Schema 强断言**:在读取粗筛数据后,必须强制比对物理列名、数据类型与预期严格一致,不一致立即阻断流水线。 |
||
| 55 | * **备注原因**:防止上游表结构发生微小异动(如字段重命名)时,静默流入千万级错误数据导致全链路重算浪费。 |
||
| 56 | * **按年拆分落盘强制豁免**:针对 S2 数据加载与 S4 标签生成模块,允许全量加载入内存,但**必须在落盘前强制按年份拆分为独立 Parquet 文件**。 |
||
| 57 | * **备注原因**:这两个模块的输入数据无法按年份谓词在存储层裁剪,但下游训练极度依赖按年懒加载,因此必须在内存中承担一次性开销完成物理拆分。 |
||
| 58 | * **S2 产出(严格 5 列)**:`fund_id`, `net_value_date`, `cumulative_net_value`, `segment_id`, `weekly_return`。 |
||
| 59 | * **S3 产出(严格 42 列宽表)**:3 列物理主键 + 38 维 Float64 原始特征 + 1 列物理哨兵 `is_feature_complete`。 |
||
| 60 | * **S4 产出(严格 4 列)**:`fund_id`, `net_value_date`, `segment_id`, `label`。 |
||
| 61 | * **Parquet 物理有序、Row Group 与压缩死锁约束**:所有特征与净值落盘前,必须按 `["net_value_date", "fund_id"]` 的字典序进行绝对物理排序;且强制设定 `row_group_size=100000`,**强制指定压缩编解码器为 `compression="zstd"` 且 `compression_level=7`**。 |
||
| 62 | * **备注原因**:下游按时间范围懒加载时,引擎依赖物理排序的 Row Group 元数据直接跳过不相关的文件块。同时,Float64 金融浮点数在 Snappy 下压缩比极差(导致巨大的无效 I/O)。从信息论实测来看,Level 7 已经能吃掉金融时序 90% 的冗余,Level 10 只是边际收益(多省 3% 磁盘),但 CPU 编码耗时呈指数级上升。在 Docker CPU Quota 限制下,Level 7 是物理极限性价比点,其解压速度远超弥补多余 I/O 的开销,是倒逼懒加载吞吐量触及物理极限的唯一选择。 |
||
| 63 | ## 3. 标签系统数学定义 (S4) |
||
| 64 | ### 3.1 物理断裂与独立计算 |
||
| 65 | 不同 `segment_id` 的数据在时间轴上绝对断裂。S4 必须严格按 `segment_id` 拆分为独立子数组,严禁跨 segment 拼接。每个子数组索引从 0 起始。 |
||
| 66 | * **S4 禁止关联特征列红线**:S4 模块绝对禁止在此阶段产出或关联任何特征列数据。 |
||
| 67 | * **备注原因**:强行关联会导致标签生成模块需要加载千万级特征宽表,引发致命 OOM 崩溃。毒样本校验已通过哨兵列延后至 S5 处理。 |
||
| 68 | ### 3.2 DCA 闭式向量化公式 |
||
| 69 | 第 $i$ 周(1-based)总绝对收益率 $R_i$ 的绝对数学定义: |
||
| 70 | $$R_i = \frac{\sum_{k=1}^{i} \frac{NV_i}{NV_k}}{i} - 1$$ |
||
| 71 | * **DCA 公式唯一合法翻译**:在工程实现中,严禁使用逐行迭代的循环逻辑计算上述公式,**必须且只能**通过“计算倒数数组的累加和,再与原数组相乘除以周数”的绝对向量化范式实现。 |
||
| 72 | * **备注原因**:千万级数据面板下,Python 级别的循环会耗时数小时。向量化范式将计算下沉至 C/Rust 底层,耗时降至秒级,是系统可实施的物理前提。为防止研发在此处自作聪明,强制规定以下为唯一合法代码骨架: |
||
| 73 | ```python |
||
| 74 | assert np.all(np.isfinite(nv_array)) and np.all(nv_array > 0), "Fatal: NV array contains inf/nan/non-positive values" |
||
| 75 | inv_nv = 1.0 / nv_array |
||
| 76 | cum_inv = np.cumsum(inv_nv) |
||
| 77 | weeks = np.arange(1, len(nv_array) + 1) |
||
| 78 | returns = (nv_array * cum_inv) / weeks - 1.0 |
||
| 79 | ``` |
||
| 80 | * **最小定投期否定声明**:严禁在标签计算中引入“最小定投期(如前 4 周不计算收益)”的业务截断逻辑。 |
||
| 81 | * **备注原因**:强制引入非对称窗口会彻底破坏数学向量化公式的闭环形态。必须严格遵循“窗口内任意位置首次达标即输出对应索引”的纯粹数学定义。 |
||
| 82 | ### 3.3 止盈判定逻辑与边界 |
||
| 83 | * 若 $R_i \ge 0.20$,判定为止盈。 |
||
| 84 | * **输出绝对边界范围强制为 `[1, 150]`**(1-based 索引)。严禁截断第 150 周的合法样本。 |
||
| 85 | * 毒样本截断的物理执行点强制延后至 S5 训练流水线,基于 `is_feature_complete` 哨兵列覆写。 |
||
| 86 | * **【标签语义不连续性消化声明】**:必须注意标签值域中 `0`(未达标)与 `1`(极早期达标)存在固有的语义断崖。严禁将其直接作为普通连续数值回归目标喂入排序模型,必须根据下游训练器的任务形态(如回归/排序)配置特殊的目标函数或权重进行消化。 |
||
| 87 | * **备注原因**:如果直接丢进排序模型,模型无法理解 0 和 1 之间的大小关系,会导致排序分全盘噪声化,模型彻底失效。 |
||
| 88 | ## 4. 特征工程刚性字典域 (S3) |
||
| 89 | 锁定 38 维。所有带窗口期的特征计算,**必须绝对遵从“分母隔离强制防线”**:显式设定 `min_periods=窗口长度`。遇 NULL 必须直接输出 NULL,严禁残缺窗口计算。 |
||
| 90 | ### 4.0 【关键架构决策 A 记录】38维否决决策 |
||
| 91 | * **冲突点**:曾计划补充 `nav_ma_26w` 与 `nav_ma_52w`(净值均线绝对值)以满足“40+维”字面要求。 |
||
| 92 | * **否决理由**:累计净值(NAV)是强路径依赖的单调递增序列。老基金与新基金的绝对均线存在巨大量纲差异。若强行执行截面标准化,模型将彻底丧失对基金存续期与净值规模的感知能力。 |
||
| 93 | * **结论**:为捍卫全量特征严格走截面的普适性铁律,主动摒弃对标准化不友好的绝对值指标。38 维已实现对定投体验的完备刻画。 |
||
| 94 | ### 4.1 物理哨兵列的生成算法动作 |
||
| 95 | * **业务动作**:在产出 38 维特征后,必须立即执行全横向非空判定生成 `is_feature_complete` 列。仅当该行数据的 38 个特征值全部不为空时,该哨兵列为真。 |
||
| 96 | * **备注原因**:这是下游截断毒样本的唯一绝对物理依据。如果不在此处固化判定,后续在训练流水线中重新做复杂的多窗口同步校验,极易因逻辑遗漏导致毒样本穿透防御网络。 |
||
| 97 | ### 4.2 跨段物理隔离红线 |
||
| 98 | * **业务动作**:所有包含滚动窗口的计算算子,其执行上下文必须且只能同时绑定“基金 ID”与“段 ID”两个维度。严禁仅按单一基金维度执行。 |
||
| 99 | * **备注原因**:如果忽略 Segment 隔离,滚动窗口会跨越被物理切断的“长缺失空隙”进行计算。这会导致原本应该断裂的剧烈波动被拉平,产生“伪低波动”特征,彻底误导模型。 |
||
| 100 | ### 4.3 特征详细字典基线 (38维刚性定义) |
||
| 101 | > **⚠️ 架构冻结声明** |
||
| 102 | > 代码实现中的特征名、标准化参数输出的特征名字段**必须与下表列出的 38 个特征名严格一字不差**。严禁擅自增删改。 |
||
| 103 | #### A. 均值回归类 (5维) |
||
| 104 | | 特征名 | 计算公式 (严格遵从规约) | 窗口期 | 业务含义与定投必要性 | |
||
| 105 | | :--- | :--- | :--- | :--- | |
||
| 106 | | `price_vs_ma_ratio_12w` | `nav / rolling_mean(nav, 12)` | 12 | **短期偏离度**。捕捉近3个月的超买超卖。 | |
||
| 107 | | `price_vs_ma_ratio_26w` | `nav / rolling_mean(nav, 26)` | 26 | **中期偏离度**。定投最核心的观测窗口。 | |
||
| 108 | | `price_vs_ma_ratio_52w` | `nav / rolling_mean(nav, 52)` | 52 | **长期偏离度**。判断历史级别的低估区域。 | |
||
| 109 | | `price_vs_ma_ratio_ewm_26w` | `nav / ewm_mean(nav, span=26)` | 26 | **指数加权偏离度**。对趋势反转灵敏度更高。 | |
||
| 110 | | `price_vs_ma_ratio_ewm_52w` | `nav / ewm_mean(nav, span=52)` | 52 | **长期指数加权偏离度**。长期趋势反转早期信号。 | |
||
| 111 | 2 | Huarui Lin | |
| 112 | |||
| 113 | 1 | Huarui Lin | #### B. 波动类 (7维) |
| 114 | 2 | Huarui Lin | |
| 115 | | 特征名 | 计算公式 (严格遵从规约) | 窗口期 | 业务含义与定投必要性 | |
||
| 116 | | :--- | :--- | :--- | :--- | |
||
| 117 | | `weekly_return` | `(nav_t / nav_{t-1}) - 1` | 无 | **微观基础动量**。作为波动与动量计算的基石。 | |
||
| 118 | | `rolling_std_12w` | `rolling_std(weekly_return, 12)` | 12 | **短期波动率**。衡量近期净值抖动程度。 | |
||
| 119 | | `rolling_std_26w` | `rolling_std(weekly_return, 26)` | 26 | **中期波动率**。定投核心波动观测期。 | |
||
| 120 | | `rolling_std_52w` | `rolling_std(weekly_return, 52)` | 52 | **长期波动率**。基金的固有波动属性标签。 | |
||
| 121 | | `downside_vol_26w` | `rolling_std(min(weekly_return, 0), 26)` | 26 | **下行波动率**。直接刻画定投被套的痛苦程度。 | |
||
| 122 | | `downside_vol_52w` | `rolling_std(min(weekly_return, 0), 52)` | 52 | **长期下行波动率**。长期风险底线的度量。 | |
||
| 123 | | `volatility_regime` | `rolling_quantile(rolling_std_12w, window=52, q=0.5)` | 12/52 | **波动率状态分位数**。利用短期波动在长期历史中的相对位置识别状态,输出天然有界 [0,1],抗单周极端脉冲长尾分布。 | |
||
| 124 | |||
| 125 | 1 | Huarui Lin | #### C. 趋势与动量类 (11维) |
| 126 | 3 | Huarui Lin | |
| 127 | 1 | Huarui Lin | > **⚠️ OLS 趋势斜率数学常量死锁声明** |
| 128 | > `trend_slope_26w` 与 `trend_slope_52w` 中的“年化”系数,**强制锁定为简单复利口径 `52`**(即直接乘以 52)。 |
||
| 129 | > **严禁使用 `sqrt(52)` 波动率口径**。原因:对数斜率本质是一阶期望而非二阶离差,套用根号年化会严重压缩特征量纲,破坏截面标准化分布,且与本规约 D 类动量特征的线性放大逻辑保持绝对一致。 |
||
| 130 | 3 | Huarui Lin | |
| 131 | 1 | Huarui Lin | | 特征名 | 计算公式 (严格遵从规约) | 窗口期 | 业务含义与定投必要性 | |
| 132 | | :--- | :--- | :--- | :--- | |
||
| 133 | | `momentum_4w` | `(nav_t / nav_{t-4}) - 1` | 4 | **月度动量**。捕捉右侧起涨点。 | |
||
| 134 | | `momentum_12w` | `(nav_t / nav_{t-12}) - 1` | 12 | **季度动量**。验证短期趋势是否确立。 | |
||
| 135 | | `momentum_26w` | `(nav_t / nav_{t-26}) - 1` | 26 | **半年动量**。中期趋势强度的核心指标。 | |
||
| 136 | | `momentum_52w` | `(nav_t / nav_{t-52})^(52/52) - 1` | 52 | **年度动量**。长期牛熊分界线。 | |
||
| 137 | | `max_drawdown_26w` | `1 - nav_t / rolling_max(nav, 26)` | 26 | **中期最大回撤**。构建低成本底仓的良机。 | |
||
| 138 | | `max_drawdown_52w` | `1 - nav_t / rolling_max(nav, 52)` | 52 | **长期最大回撤**。极端风险承受力测试。 | |
||
| 139 | | `trend_slope_26w` | `log(nav) OLS斜率 × 52` | 26 | **中期趋势斜率**。正斜率是止盈先决条件。 | |
||
| 140 | | `trend_slope_52w` | `log(nav) OLS斜率 × 52` | 52 | **长期趋势斜率**。长期盈利能力的确定性。 | |
||
| 141 | | `trend_r_squared_26w` | `OLS回归的R²` | 26 | **中期趋势稳定性**。定投的“画线”体验越好。 | |
||
| 142 | | `trend_r_squared_52w` | `OLS回归的R²` | 52 | **长期趋势稳定性**。区分“慢牛”与“剧烈震荡”。 | |
||
| 143 | | `consecutive_down_weeks` | 从当前向前连续 `weekly_return < 0` 的周数 | 无 | **连阴周数**。极佳左侧建仓信号。 | |
||
| 144 | 3 | Huarui Lin | |
| 145 | 1 | Huarui Lin | #### D. 定投特有类 (10维) |
| 146 | | 特征名 | 计算公式 (严格遵从规约) | 窗口期 | 业务含义与定投必要性 | |
||
| 147 | | :--- | :--- | :--- | :--- | |
||
| 148 | | `dca_cost_ratio_12w` | `rolling_mean(nav, 12) / nav` | 12 | **12周定投成本浮盈率**。直接刻画定投盘状态。 | |
||
| 149 | | `dca_cost_ratio_26w` | `rolling_mean(nav, 26) / nav` | 26 | **26周定投成本浮盈率**。半年度持仓盈亏体验。 | |
||
| 150 | | `dca_cost_ratio_52w` | `rolling_mean(nav, 52) / nav` | 52 | **52周定投成本浮盈率**。长线定投安全垫厚度。 | |
||
| 151 | | `dca_return_12w` | `(nav - rolling_mean(nav, 12)) / rolling_mean(nav, 12)` | 12 | **12周定投收益率**。符合基民阅读习惯。 | |
||
| 152 | | `dca_return_26w` | `(nav - rolling_mean(nav, 26)) / rolling_mean(nav, 26)` | 26 | **26周定投收益率**。 | |
||
| 153 | | `dca_return_52w` | `(nav - rolling_mean(nav, 52)) / rolling_mean(nav, 52)` | 52 | **52周定投收益率**。 | |
||
| 154 | | `dca_return_vol_26w` | `rolling_std(dca_return序列, 26)` | 26 | **定投收益波动率**。定投体验像坐过山车。 | |
||
| 155 | | `dca_return_vol_52w` | `rolling_std(dca_return序列, 52)` | 52 | **长期定投收益波动率**。 | |
||
| 156 | | `dca_win_rate_26w` | `rolling_mean(dca_return > 0, 26)` | 26 | **26周定投正收益周占比**。定投的“舒适度”指标。 | |
||
| 157 | | `dca_win_rate_52w` | `rolling_mean(dca_return > 0, 52)` | 52 | **52周定投正收益周占比**。长期持有体验指标。 | |
||
| 158 | #### E. 风险调整类 (5维) |
||
| 159 | | 特征名 | 计算公式 (严格遵从规约) | 窗口期 | 业务含义与定投必要性 | |
||
| 160 | | :--- | :--- | :--- | :--- | |
||
| 161 | | `rolling_sharpe_26w` | `(rolling_mean(weekly_ret) - weekly_rf) / rolling_std(weekly_ret) × √52` | 26 | **半年期夏普比率**。承担1单位风险的超额收益。 | |
||
| 162 | | `rolling_sharpe_52w` | `(rolling_mean(weekly_ret) - weekly_rf) / rolling_std(weekly_ret) × √52` | 52 | **一年期夏普比率**。 | |
||
| 163 | | `rolling_sortino_26w` | `(rolling_mean(weekly_ret) - weekly_rf) / rolling_std(min(weekly_ret, 0)) × √52` | 26 | **半年期索提诺比率**。Sortino比Sharpe更契合定投。 | |
||
| 164 | | `rolling_sortino_52w` | `(rolling_mean(weekly_ret) - weekly_rf) / rolling_std(min(weekly_ret, 0)) × √52` | 52 | **一年期索提诺比率**。 | |
||
| 165 | | `calmar_ratio_52w` | `((nav_t / nav_{t-52})^(52/52) - 1) / max_drawdown(52)` | 52 | **卡玛比率**。止盈效率的终极体现。 | |
||
| 166 | #### 特征工程实现红线约束重申 (v 3.2.0 强制覆写) |
||
| 167 | > **⚠️ 致命互斥清除声明** |
||
| 168 | > 研发必须以本节为准,废弃早期基线文档的冲突描述。 |
||
| 169 | 1. **存储形态红线 (绝对覆写)**:上述 38 维特征在产出文件中,**必须严格以宽表形态落盘**(Schema: 3 列物理主键 + 38 维 Float64 特征 + 1 列物理哨兵 `is_feature_complete`,共计严格 42 列)。严禁在此阶段执行转换为长表操作。 |
||
| 170 | 2. **标准化红线**:38 维特征**无条件、无特例**全部进入截面标准化与裁剪流水线。 |
||
| 171 | 3. **分母隔离红线 (绝对覆写)**:所有带窗口期的特征计算,**严禁使用任何残缺窗口容忍逻辑**。必须绝对遵从“分母隔离强制防线”,遇 NULL 必须直接输出 NULL。 |
||
| 172 | > **⚠️ OLS 趋势计算物理防线声明(数学与拓扑双重死锁) (v 3.2.4 彻底重构)** |
||
| 173 | > `trend_slope_26w`、`trend_r_squared_26w` 等 OLS 算子的计算,必须绝对遵循以下数学与物理约束: |
||
| 174 | > 1. **滑动窗口强制声明**:必须在时间序列上逐点产出(即输出序列长度与输入严格一致),严禁仅输出尾部单一标量值。 |
||
| 175 | > 2. **矩阵求解禁令**:严禁使用任何通用矩阵求逆、最小二乘矩阵分解(如 SVD、QR)算法。必须利用自变量为连续整数索引($0, 1, ..., W-1$)的数学特性,将离差平方和分母 $SS_{xx}$ 推导为 $O(1)$ 的常量公式 $W(W^2-1)/12$。 |
||
| 176 | > 3. **缺失值分母隔离防线**:滑动窗口内一旦存在缺失值,该时间点的 OLS 输出必须直接置为 NULL,严禁执行任何形式的插值或前向填充,必须交由 `is_feature_complete` 哨兵列统一拦截。 |
||
| 177 | > **备注原因**:千万级面板下,通用矩阵 API 的内部参数校验与内存分配会被放大数十万次导致函数栈雪崩;同时缺失值如果不触发 NULL 隔离,会产生“幽灵 NaN”穿透下游哨兵列防线。具体的语言级向量化实施骨架参见第二部分 9.5 节。 |
||
| 178 | ## 5. 标准化与训练拓扑 (S5) |
||
| 179 | ### 5.1 【核心决策记录】双形态参数落盘的拓扑必然性 |
||
| 180 | * **业务动作**:S5.2 标准化阶段必须产出**两份独立的参数文件**: |
||
| 181 | 1. `standardization_params_wide.parquet`:保持极小聚合宽表形态(几百行 × 38列)。 |
||
| 182 | 2. `standardization_params.parquet`:Unpivot 后的 7 列长表形态。 |
||
| 183 | * **决策原因(严禁修改)**:S5.3 训练流水线面对的是 3000 万行 × 38 列的原始特征宽表。如果只有长表参数,物理上唯一的 Join 路径是将宽表 Unpivot 成 11.4 亿行,这会在任何有限内存环境下触发瞬间 OOM 崩溃。通过保留宽表参数,S5 训练可利用底层引擎的广播机制零拷贝 Apply;而 S6 推理必须使用长表参数以支持 `asof_join` 的时间穿透。这是“空间换时间”全局原则的典型体现。 |
||
| 184 | ### 5.2 纯横截面聚合与时序对齐 |
||
| 185 | * 必须且只能按 `net_value_date` 分组计算 `mean, std, median, mad`。 |
||
| 186 | * **MAD 裁剪上下界的业务计算动作**:在极小聚合宽表阶段,必须显式基于各截面的 `median` 与 `mad` 推导出裁剪上下界(供下游 Apply)。 |
||
| 187 | * **【双形态参数的物理时序隔离声明】**:必须在极小聚合宽表阶段执行时序偏移实现 T-1 对齐,严禁在 Unpivot 后的长表上做时序偏移。但基于双文件隔离拓扑,**必须对两种形态采用绝对不同的物理对齐策略**: |
||
| 188 | 1. **训练宽表**:直接在稀疏的 `net_value_date` 索引上执行 `.shift(1)`,保证与 S5.3 的等值 Join 物理对齐。 |
||
| 189 | 2. **推理长表**:在 Unpivot 后,必须强制关联一张预生成的**绝对自然日频日历表**,将稀疏的周频参数向前填充(Forward Fill)至每日连续状态。 |
||
| 190 | * **备注原因**:金融时序跨周末/节假日存在物理断裂。若推理长表保持稀疏,周一请求拉取参数时极易因底层引擎对非交易日的处理差异产生错位。引入自然日历表补齐,彻底抹除“周五产生数据、周一拉取请求”这一真实业务场景下的时序对齐静默 Bug,且完美利用了 S5.1 决策的“双文件隔离”拓扑,互不干扰。 |
||
| 191 | * **标准化参数长表的物理排序键死锁**:用于推理服务的 7 列长表参数落盘前,**必须强制按 `["feature_name", "date"]` 的双键字典序进行绝对物理排序**。 |
||
| 192 | * **备注原因**:推理服务的核心算子强依赖 `by` 键的局部单调性。按此严格排序才能对齐底层物理要求与标准长表列名。如果乱序或写反,会导致底层在执行时间回溯匹配时发生数据穿插 Panic 或错乱关联。 |
||
| 193 | ### 5.3 生产就绪判定铁律 |
||
| 194 | * **业务动作**:在标准化参数表中,若某特征在某截面的 `std` 为 NULL 或 0.0,或 `mad` 为 NULL 或 0.0,则该截面的 `is_production_ready` 必须强制置为 `False`。 |
||
| 195 | * **备注原因**:零方差或无绝对中位差的截面属于数学上的病态分布。如果不下发此判定,推理服务在遇到这些极端情况时会触发除零崩溃,或者将毫无区分度的垃圾特征输入模型。 |
||
| 196 | ### 5.4 按折懒加载与关联防线 |
||
| 197 | * **业务动作**:由于数据按年份拆分为多个物理文件,在构建数据流时,必须强制保证文件列表严格按照时间自然序(如 2020, 2021...)排列。 |
||
| 198 | * **备注原因**:乱序读取跨年文件会破坏时间维度的连续性,导致按时间范围裁剪时无法命中底层存储的 Row Group 谓词,引发全表扫描爆炸。 |
||
| 199 | * **【训练 Join 三键死锁】**:在特征宽表与标签表进行物理关联时,必须且只能强制指定三键联合匹配:`["fund_id", "net_value_date", "segment_id"]`,严禁遗漏 `segment_id`。 |
||
| 200 | * **备注原因**:如果不强制锁死三键,极易仅按双键关联,导致同一只基金在不同 Segment 的重叠时间段内发生笛卡尔积爆炸,直接触发内存雪崩。 |
||
| 201 | ### 5.5 【核心防线】Spearman 秩相关系数矩阵的内存防爆与数学等价防线 |
||
| 202 | 在进入 5 折训练前,计算全量训练集特征间 Spearman 秩相关系数矩阵以剔除共线性(阈值 0.90)时,必须强制遵守以下纯数学算法约束: |
||
| 203 | * **横截面分块算秩**:严禁一次性对全量数据执行多列全局 `rank()` 导致内存翻倍。必须按 `net_value_date` 分批计算特征在每个截面内的秩。 |
||
| 204 | * **【分块物理标尺具象化】**:严禁按固定行数分块,**必须按自然年分块**。单截面全市场约 1.5 万只基金,38 维特征双份矩阵(原始+秩)仅占不到 10MB。按年分块峰值远低于 4GB 红线,且完美对齐 S2 产出的按年 Parquet 物理文件,触发 Row Group 级别裁剪。 |
||
| 205 | * **备注原因**:“分块”若无明确物理对齐策略,极易导致跨年边界处的数据被截断或重复计算。按年分块是唯一既满足内存限制又保证文件 I/O 零冗余的标尺。 |
||
| 206 | * **Spearman 秩计算的防漂移锁定**:在分块算秩时,**必须显式锁定“平均秩”处理方式,严禁依赖底层计算引擎的默认行为**。 |
||
| 207 | * **备注原因**:不同底层版本或不同语言对并列项的打分策略可能不同(如平均法、最大法、最小法)。如果不显式锁定,在环境迁移或底层库升级时,会导致共线性矩阵发生静默漂移,最终选出的特征子集完全不可复现。 |
||
| 208 | * **强制中心化截断**:分块算秩后,必须立即执行中心化操作,以彻底消除后续累加项量级爆炸导致的浮点精度丢失。 |
||
| 209 | * **Pearson 展开项累加**:将中心化后的全量截面秩数据视为一个大面板,利用 Pearson 协方差公式的展开项分块累加。 |
||
| 210 | * **动态有效样本量联合累加**:分母的有效样本量必须通过双变量联合非空判断同步分块独立累加,确保数学绝对严谨。 |
||
| 211 | * **【内存峰值硬限制】**:分块大小的唯一物理标尺是,必须保证单次分块计算时的内存驻留峰值绝对不超过 4GB。 |
||
| 212 | * **备注原因**:“分块”是一个模糊概念,如果没有硬限制数值,研发极可能划分过大的块依然触发 OOM,必须给一个具体的物理数值作为分块的唯一标尺。 |
||
| 213 | * **对称化强制截断**:分块累加完成后,必须强制执行对称化截断 `matrix = (matrix + matrix.T) / 2`,确保矩阵绝对对称。 |
||
| 214 | * **备注原因**:千万级数据分块累加极易引入微小浮点误差导致矩阵不对称。如果不执行此截断,可能导致后续特征剔除逻辑误判主对角线以外的值,引发共线性剔除结果偏移。 |
||
| 215 | ### 5.6 训练评估与闭环重训防线 |
||
| 216 | * **业务动作**:必须采用累积扩展窗口策略。第 $i$ 折使用前 $i \times 20\%$ 数据训练,验证集为紧接着的固定步长区块。5 折评估期间严禁再次剔除特征,仅用于指标观测。 |
||
| 217 | * **备注原因**:金融时序数据存在严重的非稳态分布漂移。累积扩展窗口能最大程度利用早期历史数据训练趋势,同时用固定步长区块模拟真实的向前推导,防止随机划分导致未来数据泄漏到训练集。 |
||
| 218 | * **【闭环重训物理动作】**:在 5 折评估仅用于指标观测完毕后,必须使用确定的特征子集对全量训练数据进行最后一次重训,并以此产出最终落盘的物理参数文件与模型文件。 |
||
| 219 | * **备注原因**:如果研发只在第 5 折评估完后就直接把模型推上线,线上推理服务拉取到的参数仍然是“预剔除阶段”基于部分数据算出来的,会导致线上特征分布与模型绝对不匹配,预测全盘崩溃。 |
||
| 220 | ### 5.7 S6 推理业务契约 |
||
| 221 | * **【核心决策记录】两步走时间锚定的业务隔离意义** |
||
| 222 | * **业务动作**:S6 推理必须分两步执行参数拉取:第一步,将请求日期与单列的 `available_param_dates` 对齐,获取有效参数日期;第二步,用有效日期去拉取具体参数。 |
||
| 223 | * **决策原因(严禁修改)**:推理服务是系统的最后一道防线。将“日期容灾验证”与“海量参数拉取”解耦,可以在遇到冷启动(请求日期早于最小参数日期)时,实现微秒级短路返回,避免引擎去触碰庞大的参数表,这是保障在线服务 P99 延迟的绝对业务铁律。 |
||
| 224 | * **推理服务业务红线**:单基金查询接口与批量推荐接口绝对隔离。批量推荐接口返回的数据结构中,用于展示归因的字段**必须强制为空字符串**。 |
||
| 225 | * **备注原因**:批量接口面对高并发必须压榨至极致算力。归因计算涉及极其昂贵的树结构遍历,如果在批量链路中触发,会直接导致服务超时熔断。 |
||
| 226 | ## 6. 回测评估基准 |
||
| 227 | * **快照策略**:采用验证集**最后一个时间截面**的单点快照进行打分分桶。 |
||
| 228 | * **分桶强制算法**:必须采用等频分位(如分为 5 桶)进行分桶评估。 |
||
| 229 | * **备注原因**:金融资产的预测分数通常呈长尾分布。等频分位能保证每个桶内有绝对足够的样本量,防止等宽分位导致尾部桶样本稀疏、回测指标失去统计学意义。 |
||
| 230 | * **零方差防线**:若该截面预测分数方差为 0(如全量预测值相同),导致无法构建分箱边界,**该截面必须强制跳过评估**,对应桶的各项指标强制输出 `NULL`,严禁抛出异常阻断回测流水线。 |
||
| 231 | * **备注原因**:模型在初期或极端行情下可能输出全量相同的保守得分。如果强行分箱会触发除零异常,阻断整个自动化回测流水线的运转。 |
||
| 232 | --- |
||
| 233 | # 第二部分:Polars/Arrow 物理实施手册 |
||
| 234 | **受众**:底层研发工程师。 |
||
| 235 | **核心原则**:性能偏执。本篇专门针对 C/Rust/Arrow 引擎的微观物理缝隙进行强制封堵,**严禁以“字面直译”替代“物理行为对齐”**。任何违反本篇的代码实现,在 Code Review 中享有一票否决权。 |
||
| 236 | ## 0. 全局工程基准原则 |
||
| 237 | > **【最高铁律】如遇到空间资源与时间资源的冲突,优先考虑使用空间换时间。** |
||
| 238 | > 例如:宁可多落盘一份宽表参数文件占用几 MB 磁盘,也严禁在 3000 万行面板上执行 Unpivot 导致 11.4 亿行内存雪崩。 |
||
| 239 | ## 1. 内存生命周期管理契约 |
||
| 240 | * **实施铁律**:释放 Arrow 内存,仅依赖 Python 引用计数的 `del df`。 |
||
| 241 | * **【封杀 gc.collect()】**:**严禁调用 `gc.collect()`**。该指令无法触达 Rust 底层的 mimalloc 分配器,不仅无法释放 Arrow 内存,反而会触发全局 Python 解释器 STW(Stop-The-World),在多线程高频切分场景下是极其严重的性能毒药。 |
||
| 242 | ## 2. 排序与落盘物理契约 |
||
| 243 | * **实施铁律**:所有 `sort()` 均采用默认参数(即 `maintain_order=False`)。 |
||
| 244 | * **【封杀 maintain_order=True】**:该参数会强制 Polars 引擎放弃最高效的 O(N) Radix Sort,退化为 O(N log N) 的 Timsort 并引入巨大比较开销。在千万级数据落盘时会导致耗时翻倍。Polars 的 Radix Sort 天然保证全局有序。 |
||
| 245 | ## 3. 分组计算与 EWM 防线 |
||
| 246 | * **实施铁律**:放弃原规约中的全局 `partition_by` 内存翻倍陷阱,改用 `df.group_by(["fund_id", "segment_id"]).map_elements(func, schema=...)` 进行流式迭代计算 EWM。 |
||
| 247 | * **【封杀无边界 partition_by 内存爆炸】(v 3.2.2 措辞降级修正)**:`df.partition_by(as_dict=True)` 会发生真实的物理内存全量拷贝。在 50GB cgroup 下面对千万级**全量或未严格测定内存波峰的拼接表**时必触发 OOM Kill,必须使用 `map_elements` 利用引擎内部的分块迭代器规避。**【局部安全池豁免】**:在已通过 Row Group 裁剪确信内存绝对安全的局部闭包内(如单年份实体表 < 3GB),允许将其作为向纯 NumPy C 层传递连续内存的桥接器使用。 |
||
| 248 | * **备注原因**:真正的灾难来源于“在不确定大小的内存空间上触发全量拷贝”,而不是“在 3GB 的安全沙箱里提取句柄”。将封杀令精确化为“无边界场景”,既保全了防线的威慑力,又为纯 NumPy 微线程拓扑留出了必要的工程呼吸空间。 |
||
| 249 | * **【EWM 计算 GIL 防线】(v 3.2.2 彻底重构)**:在 `map_elements` 的内层 Segment 闭包内,EWM 的计算**必须下沉至 `scipy.signal.lfilter` 的纯 C 递推底层执行**,严禁使用 Python 级别的 `for` 循环,严禁引入 Numba (`@njit`) 触发 LLVM 运行时副作用。 |
||
| 250 | * **备注原因**:EWM 本质是一阶 IIR 低通滤波器。Python `for` 循环会死锁 GIL 导致多核失效;Numba 虽能释放 GIL,但其依赖的 LLVM 编译链会引入巨大的 Docker 镜像膨胀,且运行时的 JIT 预热会引发不可控的内存毛刺击穿 50GB cgroup。`scipy.signal.lfilter` 作为预编译 C 扩展,在进入繁重计算前显式释放 GIL,零依赖膨胀,是唯一满足物理极限的解法。 |
||
| 251 | ## 4. NumPy 转换零拷贝契约 |
||
| 252 | * **实施铁律**:直接 `df.select("col").to_numpy()`。 |
||
| 253 | * **【封杀 rechunk() 滥用】**:单列提取默认即为单 Chunk 连续内存,无脑调用 `.rechunk()` 会触发无效的全局遍历与边界检查,纯属浪费 CPU 周期。 |
||
| 254 | ## 5. S5.3 标准化 Apply 物理防线(绝对一票否决项) |
||
| 255 | * **实施铁律**:读取 `standardization_params_wide.parquet`,与 3000 万行特征宽表通过 `net_value_date` 执行等值 Join,利用 Arrow 广播机制原地 Apply。 |
||
| 256 | * **【封杀 11.4 亿行 Unpivot】**:严禁将 3000 万行宽表 Unpivot 后与长表参数 Join。这是本手册中唯一的“死刑”级禁止项。 |
||
| 257 | * **IEEE 754 除零拦截**:在 Apply 执行 `(feature - mean) / std` 时,严禁信任隐式除零转换(Arrow 中 `Float64 / 0 = Inf` 而非 NULL)。必须使用 `pl.when` 显式拦截: |
||
| 258 | ```python |
||
| 259 | pl.when(pl.col("std") == 0).then(pl.lit(None)) |
||
| 260 | .otherwise((pl.col("feature") - pl.col("mean")) / pl.col("std")) |
||
| 261 | ``` |
||
| 262 | ## 6. S6 推理 Join 时空语义防线 |
||
| 263 | * **物理事实对齐**:Polars 的 `asof_join` 底层本质是:对于左表键 `T`,在右表中寻找满足 `Right.key <= T` 的最大键。严禁在主观语义上将“向后寻找”脑补为“寻找未来数据”而乱加 `shift` 补丁。 |
||
| 264 | * **【封杀单维排序导致的数据穿插】**: |
||
| 265 | * **第一步(时间锚定)**:左表(请求表)必须且只能按 `request_date` 升序排序,与单列的 `available_param_dates` 执行 `asof_join`。 |
||
| 266 | * **第二步(参数拉取)**:拿到 `valid_date` 后,左表必须按 `["feature_name", "valid_date"]` 升序排序,与 7 列的 `standardization_params` 执行 `asof_join`。 |
||
| 267 | * **严禁**在包含多基金或多特征的左表上,全局仅按 `request_date` 排序后直接去 Join 参数长表,这会导致 C++ 底层因 `by` 变量局部单调性被破坏而引发 Panic 或输出错乱关联。 |
||
| 268 | * **【S6 推理进程级缓存防抖】**:在推理微服务启动时,**必须将参数表单例化为进程级全局变量**,严禁在每次请求中触发 `pl.read_parquet` 或 `pl.scan_parquet`。 |
||
| 269 | * **备注原因**:如果每次请求都重新加载,高并发下会产生巨大的锁竞争与重复内存分配。参数表极小(几万行),常驻内存零压力,彻底消灭并发 I/O 瓶颈。 |
||
| 270 | ## 7. 极简防御性代码片段(可直接 Copy-Paste) |
||
| 271 | ### 7.1 S4 索引错位拦截(0-based 转 1-based) |
||
| 272 | ```python |
||
| 273 | # 物理缝隙:NumPy 天然 0-based,业务 Label 强制 1-based [1, 150] |
||
| 274 | # 致命后果:直接用 argmax 会导致第 150 周合法止盈样本丢失,触发 FATAL-2 越界 |
||
| 275 | hit_indices = np.where(returns >= 0.20)[0] |
||
| 276 | label = int(hit_indices[0] + 1) if len(hit_indices) > 0 else 0 |
||
| 277 | assert 0 <= label <= 150, f"FATAL: Label out of bounds [0, 150], got {label}" |
||
| 278 | ``` |
||
| 279 | ### 7.2 S5 等距抽样边界退化兜底 |
||
| 280 | ```python |
||
| 281 | # 物理缝隙:np.linspace 在 num > 数组长度时,静默产生重复索引 |
||
| 282 | # 致命后果:短历史基金样本被 LightGBM 极高权重重复吸收,引发不可逆过拟合 |
||
| 283 | actual_num_samples = min(max_samples, len(valid_pool)) |
||
| 284 | sampled_indices = np.linspace(0, len(valid_pool) - 1, actual_num_samples, dtype=int) |
||
| 285 | ``` |
||
| 286 | ### 7.3 S2 短缺失占位行收益置 NULL 骨架(决策点 1 物理落地) |
||
| 287 | ```python |
||
| 288 | # 决策理由:必须切断 0.0 幻影收益向下游 rolling_std 的传播。 |
||
| 289 | # 保留 0.0 会被视为合法样本,人为压低恢复期波动率,产生“伪低波动”毒特征。 |
||
| 290 | # 置 NULL 的代价是后续 12 周窗口产生 NULL,但完美触发“分母隔离防线”, |
||
| 291 | # 最终被 is_feature_complete 哨兵列捕获,在 S5 截断为 label=0。这是以局部样本量换取全局特征纯洁性的唯一正确拓扑。 |
||
| 292 | # [v 3.2.1 强制覆写] 首行 Ffill 越界防御:彻底封堵 Arrow 到 Python 的隐式类型转换缝隙 |
||
| 293 | import math |
||
| 294 | first_nv = df_segment.select(pl.col("cumulative_net_value").first()).item() |
||
| 295 | assert first_nv is not None and not (isinstance(first_nv, float) and math.isnan(first_nv)), \ |
||
| 296 | "FATAL: Segment 首行净值为空或为 NaN,触发越界防御流水线阻断。" |
||
| 297 | df_segment = df_segment.with_columns( |
||
| 298 | pl.when(pl.col("is_original") == False) |
||
| 299 | .then(pl.lit(None)) |
||
| 300 | .otherwise(pl.col("weekly_return")) |
||
| 301 | .alias("weekly_return") |
||
| 302 | ).drop("is_original") |
||
| 303 | ``` |
||
| 304 | ### 7.4 S5.3 标准化 Apply 零拷贝广播骨架(决策点 3 与全局原则物理落地) |
||
| 305 | ```python |
||
| 306 | # 决策理由与全局原则体现:【如遇到空间资源与时间资源的冲突,优先考虑使用空间换时间】 |
||
| 307 | # 宁可在 S5.2 多存一份 standardization_params_wide.parquet 占用几 MB 磁盘(空间), |
||
| 308 | # 也严禁将 3000万x38 宽表 Unpivot 成 11.4 亿行去匹配长表参数(必 OOM),以此换取毫秒级训练 Apply(时间)。 |
||
| 309 | # [v 3.2.2 新增] 内存尖峰拦截:严禁在 LazyFrame 状态下执行此 Join! |
||
| 310 | # 必须确保 features_df 是已 .collect() 的实体表,否则底层 Hash Table 动态扩容会引发不可控的内存毛刺击穿 50GB cgroup。 |
||
| 311 | assert isinstance(features_df, pl.DataFrame), "FATAL: 标准化 Join 必须在实体 DataFrame 上执行" |
||
| 312 | params_wide = pl.read_parquet("standardization_params_wide.parquet") |
||
| 313 | # 利用 Arrow 广播机制,3000万行特征表与几百行参数表按 net_value_date 等值 Join,零内存膨胀 |
||
| 314 | features_df = features_df.join(params_wide, on="net_value_date", how="left") |
||
| 315 | # 在 3000 万行面板上原地向量化 Apply |
||
| 316 | exprs = [] |
||
| 317 | for col_name in final_feature_list: |
||
| 318 | # IEEE 754 强制除零拦截(严禁直接 / std,否则产生 Inf 击穿 LightGBM) |
||
| 319 | exprs.append( |
||
| 320 | pl.when(pl.col(f"std_{col_name}") == 0) |
||
| 321 | .then(pl.lit(None)) |
||
| 322 | .otherwise( |
||
| 323 | (pl.col(col_name).clip(pl.col(f"clip_lower_{col_name}"), pl.col(f"clip_upper_{col_name}")) - pl.col(f"mean_{col_name}")) / pl.col(f"std_{col_name}") |
||
| 324 | ) |
||
| 325 | .alias(col_name) |
||
| 326 | ) |
||
| 327 | features_df = features_df.with_columns(exprs) |
||
| 328 | ``` |
||
| 329 | ### 7.5 S6 推理两步走时间锚定骨架(决策点 2 物理落地) |
||
| 330 | ```python |
||
| 331 | import polars as pl |
||
| 332 | import logging |
||
| 333 | # [v 3.2.1 新增] 进程级全局缓存防抖,消灭并发 I/O 瓶颈 |
||
| 334 | _PARAMS_LONG_CACHE = pl.read_parquet("standardization_params.parquet").sort(["feature_name", "date"]) |
||
| 335 | _AVAIL_DATES_CACHE = _PARAMS_LONG_CACHE.select("date").unique() |
||
| 336 | def inference_single_fund(request_date: pl.DataFrame) -> pl.DataFrame: |
||
| 337 | # 决策理由:将“日期容灾验证”与“海量参数拉取”解耦。 |
||
| 338 | # 遇到冷启动(请求日期早于最小参数日期)时,在第一步即可微秒级短路返回, |
||
| 339 | # 避免引擎去触碰庞大的参数表,这是保障在线微服务 P99 延迟的绝对业务铁律。 |
||
| 340 | |||
| 341 | # === 第一步:时间锚定(单维排序,严禁传入 fund_id 破坏物理布局) === |
||
| 342 | # 物理事实对齐:asof_join 底层本质是“向历史回溯寻找 <= target_date 的最大有效键” |
||
| 343 | req_anchor = request_date.sort("request_date") |
||
| 344 | req_anchor = req_anchor.join(_AVAIL_DATES_CACHE, left_on="request_date", right_on="date", how="asof") |
||
| 345 | valid_date = req_anchor["date"][0] |
||
| 346 | |||
| 347 | if valid_date is None: |
||
| 348 | logging.warning(f"Cold start triggered for request_date: {request_date[0, 'request_date']}") |
||
| 349 | # 短路返回全量 NULL 特征,严禁抛异常阻断服务 |
||
| 350 | return pl.DataFrame({ |
||
| 351 | "feature_name": final_feature_list, |
||
| 352 | "value": [None] * len(final_feature_list) |
||
| 353 | }) |
||
| 354 | |||
| 355 | # === 第二步:参数拉取(双键排序,防范数据穿插导致 C++ 底层 Panic) === |
||
| 356 | # 若为批量推荐,左表必须先展开为 (feature_name, valid_date) 的长表形态 |
||
| 357 | # 严禁在包含多特征的左表上,仅按 request_date 排序后直接去 Join 参数长表 |
||
| 358 | req_params = pl.DataFrame({ |
||
| 359 | "feature_name": final_feature_list, |
||
| 360 | "valid_date": [valid_date] * len(final_feature_list) |
||
| 361 | }) |
||
| 362 | req_params = req_params.sort(["feature_name", "valid_date"]) |
||
| 363 | |||
| 364 | # params_long 在落盘时已按 ["feature_name", "date"] 严格物理有序,直接命中极速扫描 |
||
| 365 | # [v 3.2.0 修正] 严禁在此处前置 .filter(),防止打断 Lazy 优化图触发右表全量物化导致延迟毛刺 |
||
| 366 | result = req_params.join( |
||
| 367 | _PARAMS_LONG_CACHE, |
||
| 368 | left_on=["feature_name", "valid_date"], |
||
| 369 | right_on=["feature_name", "date"], |
||
| 370 | how="asof" |
||
| 371 | ) |
||
| 372 | |||
| 373 | # [v 3.2.0 修正] 后置生产就绪过滤:在极小的左表结果上过滤,零成本防御 |
||
| 374 | result = result.filter(pl.col("is_production_ready")) |
||
| 375 | |||
| 376 | # 后续基于 result 中的 mean, std, clip_upper, clip_lower 执行单条向量化 Z-Score 即可 |
||
| 377 | return result |
||
| 378 | ``` |
||
| 379 | ## 8. 工程治理与 CI 强制红线 |
||
| 380 | * **【Pandas 零容忍】**:全链路严禁引入 Pandas。 |
||
| 381 | * **备注原因**:Pandas 基于惰性求值与复杂的 Block Manager,一旦在中间环节混入,会彻底破坏 Polars LazyFrame 的全局查询优化图,导致不可控的隐式内存物化与性能崩塌。 |
||
| 382 | * **【配置分级强制隔离】**:`MAD系数(1.4826)`、`共线性阈值(0.90)`、`无风险利率` **必须外置于 `config.yaml`**;`周频常数(52)`、`长缺失阈值` **必须硬编码为代码常量**。 |
||
| 383 | * **备注原因**:防止“魔法数字”散落于代码库各处导致不同环境训练不一致。业务超参可配,但影响底层物理数据结构的数学常量绝对不可配,以防误操作污染全量数据血缘。 |
||
| 384 | * **【EWM/OLS 金色单测与 Schema 对称差校验】**:必须强制集成于 CI 流水线。 |
||
| 385 | * **备注原因**:EWM 与 OLS 没有现成的滚动窗口 API,纯靠手工推导,极易在底层库升级时发生精度漂移;Schema 对称差校验是拦截任何研发“自作聪明”增删特征列物理主键的唯一自动化防线。 |
||
| 386 | ## 9. 极端生产环境下的微观封堵补丁(50GB cgroup 兜底) |
||
| 387 | 本节针对规约前文在千万级高频面板数据下可能暴露的微观物理缝隙,进行终极封堵。在 Code Review 中,违反本节任何一条的代码,直接打回且不计入绩效。 |
||
| 388 | ### 9.1 S2 DuckDB 粗筛周频对齐的物理防偏移 |
||
| 389 | * **实施铁律**:虽然规约 2.2 节锁定了 `1970-01-05` 作为数学基准,但在 DuckDB 转化为具体 SQL 实体时,必须防范基金净值 T+1 日历披露规则导致的物理错位。 |
||
| 390 | * **封杀 `DATE_TRUNC` 原生函数**:严禁使用 `DATE_TRUNC('week', date)` 等依赖 DuckDB 实例所在操作系统时区环境变量的内置函数。 |
||
| 391 | * **备注原因**:如果服务器时区被意外设置为 UTC+0,而业务发生在北京时间周五晚,`DATE_TRUNC` 会将周五切分为“下一周”的开端,导致该周的代表净值被硬生生踢入下一段 Segment,引发下游全链路周级别错乱。严禁使用基于 `INTERVAL` 的加减Hack(易在周三/周四产生日期错位),**必须且只能**使用纯数学整数除法结合硬编码基准日,彻底抹除时区对齐的静默风险: |
||
| 392 | * **唯一合法 SQL 实体**: |
||
| 393 | ```sql |
||
| 394 | -- 强制转换为 BIGINT 防止整数溢出,除以 7 天取整再乘回 7,精准对齐到绝对周一 |
||
| 395 | DATE '1970-01-05' + ((CAST(date AS BIGINT) - CAST(DATE '1970-01-05' AS BIGINT)) / 7) * 7 |
||
| 396 | ``` |
||
| 397 | ### 9.2 `map_elements` 的高频调度与 GIL 泄洪防御(v 3.2.2 彻底重构) |
||
| 398 | * **实施铁律**:原规约的多进程 `fund_id` 切片方案因破坏 Row Group 裁剪导致 I/O 灾难,已被彻底废弃。**严禁裸奔调用单线程 `map_elements`**,也**严禁通过 `multiprocessing.Pool` 直接传递 Polars DataFrame 对象**。 |
||
| 399 | * **唯一合法实施路径(单年份安全内存池 + 纯 NumPy 微线程)**: |
||
| 400 | 1. **按年流式隔离**:外层循环必须严格按年份读取 9.4 节产出的单年份物理文件,确保单次内存驻留绝对安全(< 3GB); |
||
| 401 | 2. **局部豁免桥接**:在单年份安全池内,利用第 3 节豁免权,执行 `df_year.partition_by(["fund_id", "segment_id"], as_dict=True)`,将数据零拷贝打散为独立连续内存的 NumPy 视图; |
||
| 402 | 3. **纯底层微线程**:利用 `concurrent.futures.ThreadPoolExecutor` 将上述字典的 values(NumPy 数组)分发至多线程。内层计算**强制约束为仅使用 `np.dot`、`np.sum` 等基础 C 算子**(严禁实例化 Polars 对象),利用 NumPy 底层短暂释放 GIL 的特性实现真并行; |
||
| 403 | 4. **内存就地释放**:年份计算完毕后直接 `sink_parquet` 落盘,通过 `del` 立即释放年份内存池,进入下一年迭代。 |
||
| 404 | * **备注原因**:通过 Python `Pool` 管道序列化 Arrow 内存会触发致命的深拷贝;基于 `fund_id` 的扫描会触发全表 I/O 灾难。退回单进程按年流式读取,在确信内存安全的沙箱内使用 `partition_by` 提取句柄,配合纯 NumPy 基础算子的微线程并发,是零新增 Docker 依赖前提下,唯一能同时打满 CPU 物理、避免 OOM 且保证 I/O 极速裁剪的工程拓扑。 |
||
| 405 | * **兜底骨架**: |
||
| 406 | ```python |
||
| 407 | from concurrent.futures import ThreadPoolExecutor |
||
| 408 | def _process_single_segment(segment_df: pl.DataFrame) -> pl.DataFrame: |
||
| 409 | # 零拷贝抽出连续内存 |
||
| 410 | nv_array = segment_df.select("cumulative_net_value").to_numpy().flatten() |
||
| 411 | # [GIL 防线] 强制约束:此处仅允许调用 np.dot/np.sum 等纯 C 算子 |
||
| 412 | # 严禁在此函数内实例化任何 pl.DataFrame 或 pl.Series |
||
| 413 | result_dict = _pure_numpy_ols_ewm_logic(nv_array) |
||
| 414 | return pl.DataFrame(result_dict) |
||
| 415 | def compute_s3_features(): |
||
| 416 | for target_year in range(2015, 2026): |
||
| 417 | # 1. 按年流式读取,内存驻留 < 3GB,绝对安全 |
||
| 418 | df_year = pl.scan_parquet(f"s2_data/s2_{target_year}.parquet").collect() |
||
| 419 | |||
| 420 | # 2. 触发第 3 节【局部安全池豁免】,作为向 C 层传递连续内存的桥接 |
||
| 421 | segments_dict = df_year.partition_by(["fund_id", "segment_id"], as_dict=True) |
||
| 422 | |||
| 423 | # 3. 纯 NumPy 微线程并发,利用底层 C 算子释放 GIL |
||
| 424 | with ThreadPoolExecutor(max_workers=16) as executor: |
||
| 425 | results = list(executor.map(_process_single_segment, segments_dict.values())) |
||
| 426 | |||
| 427 | final_year_df = pl.concat(results, rechunk=False) |
||
| 428 | final_year_df.sink_parquet(f"s3_output/s3_{target_year}.parquet") |
||
| 429 | |||
| 430 | # 4. 物理释放,让位给下一年 |
||
| 431 | del df_year, segments_dict, results |
||
| 432 | ``` |
||
| 433 | ### 9.3 `consecutive_down_weeks` 算法复杂度死锁防线 |
||
| 434 | * **实施铁律**:在 Segment 闭包内的 NumPy 层计算“从当前向前连续 `weekly_return < 0` 的周数”时,**严禁使用任何 Python 级别的 `while` 或 `for` 循环进行逐行回溯**。 |
||
| 435 | * **备注原因**:连阴周数本质是一个状态依赖的递推逻辑。如果用纯 Python 循环回溯,在遇到长达 52 周的阴跌段时,单行计算复杂度会退化为 $O(N)$,导致千万级面板的整体计算复杂度逼近 $O(N^2)$。此单一特征的耗时将抵消掉其他 37 个特征全部向量化带来的收益。必须强制使用掩码位移与累加的高效范式,将复杂度死锁在绝对的 $O(N)$。 |
||
| 436 | * **唯一合法 O(N) 骨架**: |
||
| 437 | ```python |
||
| 438 | # 在 Segment 闭包内的纯 NumPy 层执行 |
||
| 439 | arr = (weekly_returns < 0).astype(np.int32) |
||
| 440 | # 1. 标记连续下跌块的起点(当前是1,且前一个不是1,或处于索引0) |
||
| 441 | is_start = np.zeros_like(arr, dtype=bool) |
||
| 442 | is_start[1:] = (arr[1:] == 1) & (arr[:-1] == 0) |
||
| 443 | is_start[0] = (arr[0] == 1) |
||
| 444 | # 2. 计算全局下跌日累加和 |
||
| 445 | cum_arr = np.cumsum(arr) |
||
| 446 | # 3. 计算起点累加和,并向右偏移1位作为基线 |
||
| 447 | cum_start = np.cumsum(is_start) |
||
| 448 | offset = np.zeros_like(cum_arr) |
||
| 449 | offset[1:] = cum_start[:-1] |
||
| 450 | # 4. 当前连续长度 = 全局累加 - 所属块起始累加 |
||
| 451 | consecutive = cum_arr - offset |
||
| 452 | # 5. 非下跌日强制归零 |
||
| 453 | consecutive[arr == 0] = 0 |
||
| 454 | ``` |
||
| 455 | ### 9.4 S2 豁免模块落盘前的内存尖峰拦截(v 3.2.2 过度防御剥离) |
||
| 456 | * **实施铁律**:规约 2.4 节允许 S2 与 S4 模块“全量加载入内存,按年拆分落盘”。但在执行物理排序 `sort(by=["net_value_date", "fund_id"])` 时,**严禁直接对全量内存中的 DataFrame 触发 `.sort().collect()`**。 |
||
| 457 | * **备注原因**:即使规约强制使用了不维护顺序的 Radix Sort,在处理数千万行数据时,底层 C++ 引擎仍需分配巨大的额外中间 Buffer 用于桶排序。在 50GB cgroup 硬限制下,这个内存毛刺极有可能瞬间击穿阈值触发 OOM-Killer。必须将按年拆分动作前置为过滤条件,在 `LazyFrame` 层面先按年过滤,再结合 `sort(streaming=True).sink_parquet()` 进行 Out-of-core(落盘计算)循环,彻底抹平排序动作的内存波峰。 |
||
| 458 | * **兜底骨架**: |
||
| 459 | ```python |
||
| 460 | # [v 3.2.2 修正] 剥离显式提取 year 列的过度防御。 |
||
| 461 | # 现代 Polars 对 Date32 的 .dt.year() 谓词下推已极其成熟稳定,直接使用可避免全量物化。 |
||
| 462 | for target_year in range(2015, 2026): |
||
| 463 | # 此时的 filter 能 100% 命中 Parquet 的统计信息,实现极速裁剪 |
||
| 464 | (pl.scan_parquet("raw_s2_data.parquet") |
||
| 465 | .filter(pl.col("net_value_date").dt.year() == target_year) |
||
| 466 | .sort(by=["net_value_date", "fund_id"], maintain_order=False) |
||
| 467 | .sink_parquet(f"output/s2_{target_year}.parquet", |
||
| 468 | row_group_size=100000, |
||
| 469 | compression="zstd", |
||
| 470 | compression_level=7)) |
||
| 471 | ``` |
||
| 472 | ### 9.5 OLS 滚动算子函数栈击穿与 NaN 防线(v 3.2.4 彻底重构) |
||
| 473 | * **实施铁律**:在微线程并发闭包内实现 OLS 趋势计算时,**严禁调用 `np.linalg.lstsq`、`np.linalg.inv` 等任何通用线性代数 API,严禁在闭包内实例化任何二维矩阵**。 |
||
| 474 | * **唯一合法实施路径(滑动窗口累加和展开式 + NaN 绝对拦截)**:必须利用 OLS 自变量为连续整数索引的特性,将矩阵求解手工展开为纯标量运算,**强制约束整个计算过程仅允许调用 `np.sum` 与 `np.dot` 两个基础 C 算子**。必须产出与输入序列严格等长的滚动结果,且**绝对禁止 Python `for` 循环**。 |
||
| 475 | * **备注原因**:千万级面板切分为数十万个微型 Segment 的语境下,通用线性代数 API 内部的参数校验、临时矩阵内存分配与 Python 包装器开销会被放大数十万次,产生灾难性的函数栈雪崩。只有利用连续整数特性将分母 $SS_{xx}$ 转化为 $O(1)$ 数学公式,分子转为 `np.dot`,才能将函数栈开销降至绝对的物理零点。同时,必须通过零拷贝滑动窗口视图与忽略 NaN 的累加和(`nancumsum`)配合,严格锁定 $O(N)$ 复杂度并完美拦截第一部分 4.3 节规定的 NaN 分母隔离防线。 |
||
| 476 | * **唯一合法 OLS 滚动向量化骨架**: |
||
| 477 | ```python |
||
| 478 | import numpy as np |
||
| 479 | from numpy.lib.stride_tricks import sliding_window_view |
||
| 480 | # 硬编码数学常量 |
||
| 481 | ANNUALIZATION_FACTOR = 52 |
||
| 482 | def _rolling_ols_logic(log_nv_array: np.ndarray, window: int) -> dict: |
||
| 483 | """ |
||
| 484 | 极致性能版滚动 OLS。输入可以是包含 NaN 的纯净 log(nav) 序列。 |
||
| 485 | 输出严格等长序列,遇 NaN 直接输出 None。 |
||
| 486 | 严禁在此函数内引入任何 np.linalg、矩阵实例化或 Python for 循环! |
||
| 487 | """ |
||
| 488 | n = len(log_nv_array) |
||
| 489 | slope_arr = np.full(n, np.nan, dtype=np.float64) |
||
| 490 | r2_arr = np.full(n, np.nan, dtype=np.float64) |
||
| 491 | |||
| 492 | if n < window: |
||
| 493 | return {"trend_slope": slope_arr, "trend_r_squared": r2_arr} |
||
| 494 | # 1. 构造自变量滑动视图 [O(1) 零拷贝] |
||
| 495 | x_centered = np.arange(window, dtype=np.float64) - ((window - 1) / 2.0) |
||
| 496 | |||
| 497 | # 2. 利用连续整数拓扑,O(1) 计算分母 SS_xx |
||
| 498 | ss_xx = (window * (window**2 - 1)) / 12.0 |
||
| 499 | |||
| 500 | # 3. 构造因变量滑动视图 [O(1) 零拷贝] |
||
| 501 | y_windows = sliding_window_view(log_nv_array, window_shape=window) |
||
| 502 | |||
| 503 | # 4. [核心防线] 计算每个窗口内的 NaN 数量,触发分母隔离 |
||
| 504 | nan_counts = np.isnan(y_windows).sum(axis=1) |
||
| 505 | valid_mask = (nan_counts == 0) |
||
| 506 | |||
| 507 | # 提取无 NaN 的纯净窗口切片 [底层 C 内存连续,无拷贝] |
||
| 508 | valid_y = y_windows[valid_mask] |
||
| 509 | |||
| 510 | if valid_y.shape[0] > 0: |
||
| 511 | # 5. 向量化计算 y_mean 与 y_centered [纯 C 算子,短暂释放 GIL] |
||
| 512 | y_means = valid_y.mean(axis=1, keepdims=True) |
||
| 513 | y_centered = valid_y - y_means |
||
| 514 | |||
| 515 | # 6. 纯 np.dot 计算分子 SS_xy [直接命中 SIMD 指令] |
||
| 516 | ss_xy = np.dot(y_centered, x_centered) |
||
| 517 | |||
| 518 | # 7. 计算斜率并年化 |
||
| 519 | slopes = (ss_xy / ss_xx) * ANNUALIZATION_FACTOR |
||
| 520 | |||
| 521 | # 8. 纯 np.sum 计算 R² |
||
| 522 | ss_tot = np.sum(y_centered ** 2, axis=1) |
||
| 523 | y_pred = y_means + slopes[:, np.newaxis] / ANNUALIZATION_FACTOR * x_centered |
||
| 524 | ss_res = np.sum((valid_y - y_pred) ** 2, axis=1) |
||
| 525 | |||
| 526 | # 极端边界防御:总方差为 0 时输出 None,防止 0/0 产生 NaN |
||
| 527 | r_squared = np.where(ss_tot > 0, 1.0 - (ss_res / ss_tot), np.nan) |
||
| 528 | |||
| 529 | # 9. 将合法结果回填至绝对等长数组 |
||
| 530 | slope_arr[valid_mask] = slopes |
||
| 531 | r2_arr[valid_mask] = r_squared |
||
| 532 | # NaN 会自动被下游 Polars 视为 NULL,完美触发 is_feature_complete 哨兵拦截 |
||
| 533 | return {"trend_slope": slope_arr, "trend_r_squared": r2_arr} |
||
| 534 | ``` |
||
| 535 | ### 9.6 Spearman 分块算秩 I/O 放大的环境自适应开关(v 3.2.2 强制补充) |
||
| 536 | * **实施铁律**:针对第一部分 5.5 节规定的“按自然年分块算秩”策略,必须在工程配置中增设自适应开关 `spearman.load_mode`,以应对不同部署拓扑下的 I/O 放大代价。 |
||
| 537 | * **业务动作**: |
||
| 538 | * `"yearly_chunked"`(默认值):严格按年分块,内存峰值 < 4GB。适用于 50GB cgroup 生产容器。 |
||
| 539 | * `"full_load"`(大内存直通模式):绕过按年分块,一次性加载全量特征宽表。严禁在 50GB 容器中启用。适用于 128GB+ 专用计算节点。但以下数学约束**仍然绝对生效**:按截面算秩、显式锁定平均秩、强制中心化、动态有效样本量联合累加、对称化强制截断。 |
||
| 540 | * **备注原因**:按年分块虽绝对安全,但在 10 年跨度下会将同一批 Parquet 文件重复读取约 10 次,导致 100GB+ 的纯磁盘 I/O 放大。在大内存节点强制走分块属于用安全的镣铐锁死合理性能。通过默认值锁死安全模式,杜绝生产环境误开全量模式。 |
||
| 541 | * **兜底骨架**: |
||
| 542 | ```python |
||
| 543 | import polars as pl |
||
| 544 | import yaml |
||
| 545 | # 从 config.yaml 加载,严禁硬编码 |
||
| 546 | with open("config.yaml", "r") as f: |
||
| 547 | _cfg = yaml.safe_load(f) |
||
| 548 | |||
| 549 | _SPEARMAN_MODE = _cfg["spearman"]["load_mode"] |
||
| 550 | assert _SPEARMAN_MODE in ("yearly_chunked", "full_load"), \ |
||
| 551 | f"FATAL: spearman.load_mode 必须为 'yearly_chunked' 或 'full_load', got {_SPEARMAN_MODE}" |
||
| 552 | def compute_spearman_matrix(feature_cols: list[str]) -> np.ndarray: |
||
| 553 | """ |
||
| 554 | Spearman 秩相关系数矩阵。38x38 对称矩阵。 |
||
| 555 | 数学等价性由中心化 + Pearson 展开项累加保证。 |
||
| 556 | """ |
||
| 557 | n_features = len(feature_cols) |
||
| 558 | # Pearson 展开项累加器(跨分块持久化) |
||
| 559 | sum_xy = np.zeros((n_features, n_features), dtype=np.float64) |
||
| 560 | sum_x = np.zeros(n_features, dtype=np.float64) |
||
| 561 | sum_xx = np.zeros(n_features, dtype=np.float64) |
||
| 562 | count_xy = np.zeros((n_features, n_features), dtype=np.float64) |
||
| 563 | |||
| 564 | if _SPEARMAN_MODE == "full_load": |
||
| 565 | # === 大内存直通模式:一次性全量加载 === |
||
| 566 | all_years = [] |
||
| 567 | for yr in range(2015, 2026): |
||
| 568 | all_years.append(pl.scan_parquet(f"s3_output/s3_{yr}.parquet")) |
||
| 569 | df_full = pl.concat(all_years).collect() |
||
| 570 | _accumulate_pearson_terms(df_full, feature_cols, sum_xy, sum_x, sum_xx, count_xy) |
||
| 571 | del df_full # 立即释放,避免与下游训练流水线内存叠加 |
||
| 572 | else: |
||
| 573 | # === 默认安全模式:按年流式累加 === |
||
| 574 | for yr in range(2015, 2026): |
||
| 575 | df_year = pl.scan_parquet(f"s3_output/s3_{yr}.parquet").collect() |
||
| 576 | _accumulate_pearson_terms(df_year, feature_cols, sum_xy, sum_x, sum_xx, count_xy) |
||
| 577 | del df_year # 年份级释放,保证内存峰值 < 4GB |
||
| 578 | |||
| 579 | # === 最终组装(两种模式共享,数学绝对等价) === |
||
| 580 | matrix = _assemble_pearson_from_terms(sum_xy, sum_x, sum_xx, count_xy) |
||
| 581 | |||
| 582 | # 对称化强制截断,封堵浮点误差 |
||
| 583 | matrix = (matrix + matrix.T) / 2.0 |
||
| 584 | return matrix |
||
| 585 | def _accumulate_pearson_terms(df, feature_cols, sum_xy, sum_x, sum_xx, count_xy): |
||
| 586 | """按截面算秩 -> 中心化 -> 累加 Pearson 展开项(核心数学内核)""" |
||
| 587 | from scipy.stats import rankdata |
||
| 588 | |||
| 589 | for _, group_df in df.group_by("net_value_date"): |
||
| 590 | n = len(group_df) |
||
| 591 | if n < 2: |
||
| 592 | continue |
||
| 593 | |||
| 594 | raw_matrix = group_df.select(feature_cols).to_numpy() |
||
| 595 | |||
| 596 | # 显式锁定“平均秩” |
||
| 597 | rank_matrix = np.column_stack([rankdata(raw_matrix[:, i], method="average") for i in range(raw_matrix.shape[1])]) |
||
| 598 | rank_centered = rank_matrix - rank_matrix.mean(axis=0) |
||
| 599 | |||
| 600 | for i in range(len(feature_cols)): |
||
| 601 | col_i = rank_centered[:, i] |
||
| 602 | valid_i = ~np.isnan(col_i) |
||
| 603 | n_valid_i = valid_i.sum() |
||
| 604 | if n_valid_i < 2: |
||
| 605 | continue |
||
| 606 | |||
| 607 | sum_x[i] += np.sum(col_i[valid_i]) |
||
| 608 | sum_xx[i] += np.dot(col_i[valid_i], col_i[valid_i]) |
||
| 609 | |||
| 610 | for j in range(i, len(feature_cols)): |
||
| 611 | col_j = rank_centered[:, j] |
||
| 612 | valid_mask = valid_i & (~np.isnan(col_j)) |
||
| 613 | cnt = valid_mask.sum() |
||
| 614 | if cnt < 2: |
||
| 615 | continue |
||
| 616 | |||
| 617 | sum_xy[i, j] += np.dot(col_i[valid_mask], col_j[valid_mask]) |
||
| 618 | count_xy[i, j] += cnt |
||
| 619 | |||
| 620 | # 利用对称性,将上三角复制到下三角 |
||
| 621 | idx = np.triu_indices(len(feature_cols), k=1) |
||
| 622 | sum_xy[idx[1], idx[0]] = sum_xy[idx] |
||
| 623 | count_xy[idx[1], idx[0]] = count_xy[idx] |
||
| 624 | def _assemble_pearson_from_terms(sum_xy, sum_x, sum_xx, count_xy): |
||
| 625 | n = len(sum_x) |
||
| 626 | matrix = np.zeros((n, n), dtype=np.float64) |
||
| 627 | for i in range(n): |
||
| 628 | for j in range(n): |
||
| 629 | cnt = count_xy[i, j] |
||
| 630 | if cnt < 2: |
||
| 631 | matrix[i, j] = 0.0; continue |
||
| 632 | numerator = sum_xy[i, j] - (sum_x[i] * sum_x[j]) / cnt |
||
| 633 | denominator = np.sqrt((sum_xx[i] - sum_x[i]**2 / cnt) * (sum_xx[j] - sum_x[j]**2 / cnt)) |
||
| 634 | matrix[i, j] = numerator / denominator if denominator > 0 else 0.0 |
||
| 635 | return matrix |
||
| 636 | ``` |