方案 2 框架实现 · run_factor_pipeline_regime

顺序流水线(全 Rust)· 状态进入计算过程
为什么需要新框架:第 t 天的计算参数依赖 [t−W, t−1] 的状态标量表(跨日依赖)——现有 cross_section pipeline 是 per-date 独立并行(worker 无共享状态),无法支持。 方案 2 改为主进程内严格按日期顺序推进:每算完一天就存当天标量,下一天开算前先读最近 W 天标量 → 聚类定状态 → 状态化计算。 并行度从"日期维度"转移到"日内维度"(rayon 双层并行)。
每日流水线(顺序循环,5 步)
① 读标量表
state_store 读
[t−W, t−1](14×W KB 级)
② 聚类定状态
fit_regime:雅可比 PCA
+ k=4 聚类 + 修正阈值
+ 卖压软编码(毫秒级)
③ 状态化计算
derive_compute_params →
窗口/阈值缩放生效
983 因子 + 14 状态标量
④ 写当日标量
14 个 f64 → {date}.bin
为 t+1 备料
⑤ 写因子 colblk
ShardedBackupSink
8 shard 并行 + 投影
预热期(窗口不足 W):只做 ①③④(存标量),不写因子——前 59 天不算因子,第 60 天起输出(用户定案)
⚡ 速度优化(不能跨日期并行,框架内优化)
① 单遍读内存复用:全市场逐笔一次读入内存(read_all_market),直方图与 per-stock 共用 —— 磁盘 IO 减半
② 预读流水线:独立读线程提前读未来 2 天(crossbeam bounded(2) 队列,内存可控),读盘与计算重叠 → 单日 ≈ max(读, 算)
③ 日内 rayon 双层并行:读全市场 + per-stock 各一层 par_iter
④ 写盘解耦:8 shard 无锁并行写(复用现有 ShardedBackupSink
🔧 与现有 pipeline 保持一致的细节
组件复用:colblk 存储/分片/投影(finish_and_project)、TaskResult 格式、_completed_dates 断点续算、17 特征/110 统计量核心 —— 全部与方案 1 同一套
状态引擎全 Rustfit_regime 雅可比特征分解(确定性,无随机)+ 固定种子 k-means
断点续算:因子日写 _completed_dates,重跑只处理未完成日期;预热期不标记(保证窗口充足后因子必写)
Python 调用(与 run_factor_pipeline_cross_section 同风格,单入口)
rp.run_factor_pipeline_regime(
    pipeline="peaks_regime",
    tasks=rp.td.get_range(20150105, 最新),   # 日期升序,顺序执行
    n_jobs=192,
    expected_result_length=997,
    trading_days=list(rp.td.trading_days),
    state_store_dir="/hdd/.../state_store_peaks_regime",  # 方案2 自己的标量表存储
    window=59,                                            # 滚动窗口 W
    store_dir="/hdd/.../factor_store_peaks_regime",       # colblk 输出(与方案1 完全独立)
    store_factor_names=rp.py_peaks_names(),
)