全流程 Rust 化 · Python 仅调用

读数据 → 算因子 → 写备份 → 滚动重估 → 组合,全部 Rust;Python 4 行调用
997
peaks_metrics 输出列/股/日(983 因子 + 14 状态标量伪列)
3 个
Rust 组件:peaks_metrics + peaks_regime + 补丁 A
0 行
Python 算法逻辑——只有 pipeline / 函数调用
① peaks_metrics.rs(worker 内)
rayon 读全市场逐笔 → 5 套高峰 + 110 统计量
→ 983 因子 + 14 状态标量(当日副产品)
v1 读盘 / v2 传数据
② colblk 备份写入(补丁 A)
8 shard × 997 列 → 投影
append_batch 已投影时自动降级再追加
唯一引擎改动
③ peaks_regime.rs · fit(主进程)
Rust 读 colblk 14 标量列 → 状态标量表
滚动 PCA + k=4 聚类 + 阈值 + 卖压软编码
→ 每日 params(毫秒级)
④ peaks_regime.rs · compose(主进程)
读 983 列 → 符号表 / 族乘数 / 黑名单 17
/ 门控 → 变换 → 写新 colblk + 投影
(按日期流式,内存可控)
跨日聚类不能进并行 worker(per-date 独立任务)→ 做成 Rust 函数,主进程调用 聚类输入 KB 级标量表,不读 Level2 tail_pipeline_engine 已支持 names 子集 回测 Python 调 dw(引擎本身即 Python 调度 + Rust 计算)
Python 调用文件(自上而下,无算法逻辑)
import rust_pyfunc as rp
import design_whatever as dw

names = rp.py_peaks_names()                                # 997 个因子名(Rust 单源)
rp.run_factor_pipeline_cross_section(                      # ① 读→算→写备份(全 Rust)
    pipeline="peaks", tasks=tasks, n_jobs=400,
    expected_result_length=997, update_mode=..., store_dir=store,
    store_factor_names=names, trading_days=list(rp.td.trading_days))
params = rp.peaks_regime_fit(store)                        # ③ 滚动重估(全 Rust)
rp.peaks_compose(store, out_store, arm="B", params=params) # ④ 组合→写备份(全 Rust)
dw.tail_pipeline_engine(out_store, names=names[:983], ...) # 回测(排除状态标量列)