0002-Sprint5实施步骤概要 » 历史记录 » 版本 1
Huarui Lin, 2026-04-25 22:58
| 1 | 1 | Huarui Lin | |
|---|---|---|---|
| 2 | ### Sprint 5 执行步骤与交付契约 |
||
| 3 | #### Step 4.0: 前置准备与配置收敛 |
||
| 4 | **执行动作**: |
||
| 5 | 1. 通过 `core.config_loader` 加载更新后的 `config.yaml`,提取 `s5.data_start_year` (2015) 与 `s5.data_end_year` (2026)。 |
||
| 6 | 2. 扫描 `data/processed/s3_*.parquet` 与 `data/processed/s4_*.parquet`,校验年份文件集合的对称性(必须完全覆盖配置指定的年份区间),缺失则阻断。 |
||
| 7 | **交付内容**: |
||
| 8 | * 更新后的 `configs/config.yaml` 及 `src/core/config_loader.py`(新增年份与内存阈值强校验字段)。 |
||
| 9 | * 更新后的 `src/core/constants.py`(新增 `STREAM_HARD_MEMORY_CEILING_GB = 15.0` 硬编码常量)。 |
||
| 10 | #### Step 4.1: 纯横截面聚合(极小宽表生成) |
||
| 11 | **执行动作**: |
||
| 12 | 1. 采用“按年 Lazy Scan 聚合”路径,循环扫描 S3 年份文件。 |
||
| 13 | 2. 在 LazyFrame 状态下执行 `group_by("net_value_date").agg()`,仅计算 38 维特征的 `mean, std, median, mad`。 |
||
| 14 | 3. 将各年产出的极小 DataFrame(每个截面仅 1 行)零成本 `concat` 为完整的聚合宽表。 |
||
| 15 | 4. 基于宽表向量化推导 MAD 裁剪上下界,并执行生产就绪判定(`std==0 | mad==0 -> is_production_ready=False`)。 |
||
| 16 | **交付内容**: |
||
| 17 | * 内存中的聚合极小宽表 DataFrame(Schema: `net_value_date` + 38维统计列 + 上下界列 + 哨兵列)。 |
||
| 18 | #### Step 4.2: 双形态参数物理时序对齐与落盘 |
||
| 19 | **执行动作**: |
||
| 20 | 1. **路径 A(训练宽表)**:基于 Step 4.1 产出的宽表,在稀疏的 `net_value_date` 索引上执行 `.shift(1)` 实现 T-1 对齐。调用 `core.physical_storage_contract`(使用默认排序键)落盘为 `standardization_params_wide.parquet`。 |
||
| 21 | 2. **路径 B(推理长表)**: |
||
| 22 | - 将宽表 Unpivot 为长表形态。 |
||
| 23 | - 扫描 S3 全量数据获取全局最小/最大日期极值,动态生成绝对自然日历表。 |
||
| 24 | - 将长表与自然日历表 Left Join,执行 `forward_fill` 补齐非交易日后,后置过滤掉未来日期(仅保留 `<= max_date` 的数据)。 |
||
| 25 | - 调用升级后的 `core.physical_storage_contract`,**显式传入 `override_sort_keys=["feature_name", "date"]`**,落盘为 `standardization_params.parquet`。 |
||
| 26 | **交付内容**: |
||
| 27 | * `artifacts/standardization_params_wide.parquet` |
||
| 28 | * `artifacts/standardization_params.parquet` |
||
| 29 | #### Step 4.3: 按年流式标准化 Apply |
||
| 30 | **执行动作**: |
||
| 31 | 1. 循环读取 S3 年份文件并 `.collect()` 为实体表。 |
||
| 32 | 2. 与 `standardization_params_wide.parquet` 通过 `net_value_date` 执行等值 Join。 |
||
| 33 | 3. 遍历 38 个特征列,构建 `pl.when(std==0).then(None).otherwise(clip -> sub -> div)` 的 IEEE 754 拦截表达式,原地 Apply。 |
||
| 34 | 4. 调用 `core.physical_storage_contract`(使用默认排序键),按年落盘标准化后的特征。 |
||
| 35 | **交付内容**: |
||
| 36 | * `data/processed/s5_standardized_{year}.parquet`(严格 42 列,与 S3 Schema 结构一致,仅数值被标准化)。 |
||
| 37 | #### Step 4.4: 标签关联与毒样本截断 |
||
| 38 | **执行动作**: |
||
| 39 | 1. 循环年份,按自然序读取 S4 标签表与 Step 4.3 产出的标准化表。 |
||
| 40 | 2. 执行**三键死锁 Join**:`["fund_id", "net_value_date", "segment_id"]`。 |
||
| 41 | 3. 哨兵覆写:使用 `pl.when(is_feature_complete == False).then(pl.lit(0)).otherwise(pl.col("label"))` 截断毒样本。 |
||
| 42 | 4. 按年落盘最终的纯净训练集。 |
||
| 43 | **交付内容**: |
||
| 44 | * `data/processed/s5_train_joined_{year}.parquet`(Schema: 3主键 + 38维标准化特征 + 1标签 + 1哨兵)。 |
||
| 45 | #### Step 4.5: Spearman 秩相关内存防爆防线 |
||
| 46 | **执行动作**: |
||
| 47 | 1. 严格复用 `core.stream_engine`(此时受 15GB 硬天花板保护),按年流式加载 Step 4.4 产出的训练集。 |
||
| 48 | 2. 在每个年份的闭包内,按 `net_value_date` 分组,调用 `scipy.stats.rankdata(method="average")` 显式锁定平均秩算秩。 |
||
| 49 | 3. 算秩后立即执行中心化,利用 Pearson 展开项跨年累加 `sum_xy, sum_x, sum_xx, count_xy`。 |
||
| 50 | 4. 循环结束后,组装矩阵并执行 `(matrix + matrix.T) / 2.0` 强制对称化截断。 |
||
| 51 | 5. 基于 `config.yaml` 中的 `spearman.threshold` (0.90) 剔除共线性特征。 |
||
| 52 | **交付内容**: |
||
| 53 | * `artifacts/final_feature_list.json`(降维后的最终特征子集列表)。 |
||
| 54 | --- |
||
| 55 | ### 🛑 必须执行的文档/代码修改清单 |
||
| 56 | 为了支撑上述步骤,您必须在 Step 4.0 之前完成以下文件的修改: |
||
| 57 | #### 1. `configs/config.yaml` (修改) |
||
| 58 | 新增以下节点: |
||
| 59 | ```yaml |
||
| 60 | # [S5 数据流] 动态时间窗口边界 |
||
| 61 | s5: |
||
| 62 | data_start_year: 2015 |
||
| 63 | data_end_year: 2026 |
||
| 64 | # [Core 引擎] 流式安全池弹性断言 (严禁在 50GB 生产环境设为超过 10.0) |
||
| 65 | engine: |
||
| 66 | stream_max_memory_peak_gb: 3.0 |
||
| 67 | ``` |
||
| 68 | #### 2. `src/core/config_loader.py` (修改) |
||
| 69 | 在 `PipelineConfig` 中新增对应的 Pydantic 字段,以实现启动强校验: |
||
| 70 | ```python |
||
| 71 | class S5Config(BaseModel): |
||
| 72 | data_start_year: int |
||
| 73 | data_end_year: int |
||
| 74 | class EngineConfig(BaseModel): |
||
| 75 | stream_max_memory_peak_gb: float |
||
| 76 | class PipelineConfig(BaseModel): |
||
| 77 | # ... 原有字段 ... |
||
| 78 | s5: S5Config = Field(...) |
||
| 79 | engine: EngineConfig = Field(...) |
||
| 80 | ``` |
||
| 81 | #### 3. `src/core/constants.py` (修改) |
||
| 82 | 新增硬编码天花板常量: |
||
| 83 | ```python |
||
| 84 | STREAM_HARD_MEMORY_CEILING_GB = 15.0 |
||
| 85 | ``` |
||
| 86 | #### 4. `docs/baseline/2-dry_module_list.md` (修改) |
||
| 87 | 按照会议纪要决议,精准修改以下 3 个模块的描述: |
||
| 88 | * **`core.standardization_math`**:将“被依赖方”修改为“S6 推理服务(注:S5 训练阶段因 3000万行宽表 OOM 红线已豁免复用)”。将“复用意义”修改为“S6 推理作为该内核的唯一复用方...”。 |
||
| 89 | * **`core.physical_storage_contract`**:在“强制入参”中增加 `override_sort_keys: Optional[List[str]] = None`。在“内部死锁逻辑”中增加“若传入 `override_sort_keys`,则替换默认的物理排序键,否则强制注入默认键”。 |
||
| 90 | * **`core.stream_engine`**:将“强制断言内存峰值 < 3GB”修改为“强制断言内存峰值 < min(config.engine.stream_max_memory_peak_gb, 常量.STREAM_HARD_MEMORY_CEILING_GB)”。 |
||
| 91 | #### 5. `docs/baseline/3-steps.md` (修改) |
||
| 92 | 精准替换以下 4 处描述: |
||
| 93 | * **Step 2.0 第 2 点**:“单次内存驻留强制锁定在 < 3GB 绝对安全红线内” 替换为 “单次内存驻留受 `core.stream_engine` 保护,严格执行‘配置化阈值与 15GB 硬天花板取小’的弹性断言”。 |
||
| 94 | * **Step 4.1 第 3 点**:“按 `["feature_name", "date"]` 绝对物理排序落盘” 替换为 “调用升级后的 `core.physical_storage_contract`(显式传入 `override_sort_keys=["feature_name", "date"]`)按双键绝对物理排序落盘”。 |
||
| 95 | * **Step 4.4 第 5 点**:“输出 `final_feature_list`” 替换为 “将最终特征子集物理落盘为 `final_feature_list.json`”。 |
||
| 96 | * **Step 6.3 第 2 点**:“复用 S5 内核:调用 `core.standardization_math.apply_zscore`” 替换为 “调用 `core.standardization_math.apply_zscore`(注:S5 训练阶段因内存红线已豁免复用,S6 推理为该内核的唯一复用方)”。 |