行为
03-全链路企业级算力防线蓝图 » 历史记录 » 修订 13
« 上一页 |
修订 13/14
(差异)
| 下一页 »
Huarui Lin, 2026-04-24 21:27
# 全链路企业级算力防线蓝图 (严格 0-7 阶段版)¶
第 0 阶段:项目脚手架与 CI 基线¶
本阶段是全链路的“宪法”,不写任何业务逻辑,仅确立物理约束与拦截规则。
Step 0.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 拓扑破坏。 -
依赖锁死:在
pyproject.toml中精确锁定polars,duckdb,lightgbm,numpy,scipy,fastapi,pydantic的大版本号,彻底封杀上游破坏性更新引发的静默精度漂移。【依赖纯净化红线】 严禁引入实验追踪库(如 mlflow、weights&biases)与超参搜索库(如 optuna、ray[tune])。规约 S5.2 已在拓扑层面彻底否决超参搜索,此类依赖仅会污染 Docker 镜像层与增加无用编译时间。
Step 0.2: 配置中心统一收敛¶
- 将 D-1 基金类型白名单、MAD 裁剪系数
1.4826与宽容度倍数5.0、Spearman 阈值0.90全部沉淀于config.yaml。 - 建立
core/config_loader.py,启动时做 Pydantic 强校验,缺失配置直接sys.exit(1)阻断。
Step 0.3: CI/CD 红线拦截器部署 (强制门禁)¶
在 Git 流水线与本地 Pre-commit 中植入四大静态/动态拦截器:
-
Pandas 零容忍:基于现有 AST 语法树扫描器
check_forbidden_imports.py,发现import pandas直接报错阻断。 -
Schema 对称差拦截 (已解封内聚):通过 Pytest 注册
schema_drift标记,在 S2/S3/S4 物理落盘后立刻触发强校验,比对产出文件的列名集合与硬编码契约的对称差,差集非空直接sys.exit(1)。 -
底层算子金色单测 (已解封内聚):通过 Pytest 注册
golden_precision标记,在 Step 2.4 OLS/EWM 算子实现时生成极小确定性.npy基线,后续每次跑单测强制比对,浮点误差越过1e-7直接阻断。 -
工程治理正则扫描:新增纯 Python 脚本
check_engineering_governance.py。全局拦截违规 API 与魔法数字硬编码。内置注释过滤逻辑,防止误杀文档说明。
第 1 阶段:S2 - 数据血缘与清洗拓扑¶
面对 2800 万行原始净值,走“DuckDB 粗筛压榨 -> Polars 流式微操 -> Out-of-core 物理落盘”降维路径。
Step 1.1: DuckDB 纯粗筛算子防火墙¶
-
底层转码:Polars 读取原始 CSV,强制指定类型(
fund_id为Categorical,日期为Date32),落盘为原生 Parquet 数据湖。严禁 DuckDB 直接对抗 CSV。 -
D-1 全局物理下推:在 DuckDB 内基于配置白名单过滤基础信息,提取合法
fund_id集合。 -
周频无状态聚合:利用唯一合法的 SQL 实体(基于
1970-01-05基准的纯数学整数除法对齐周一),Inner Join 后按周取组内max(date)对应净值。粗筛产物导出为极小的s2_raw_weekly.parquet,关停 DuckDB。 -
统筹编排器物理接入:本步骤严禁作为孤立脚本直接被执行,必须作为标准入口参数被
scripts/run_offline_pipeline.py极简统筹器按拓扑顺序显式调用。统筹器仅负责全局路径解析与上游异常阻断,严禁在此处包含任何业务计算逻辑,保证各 Step 模块的绝对逻辑闭包与可独立测试性。 -
结构化性能遥测挂载:本步骤的主执行函数必须强制挂载
core.telemetry.profile_step("S1.1")装饰器。执行完毕后,必须向stderr独立输出一条包含duration_sec(绝对耗时)与peak_rss_mb(内存波峰)的标准 JSON 字符串。严禁高频采样 CPU,彻底隔离正常业务数据流,为全链路性能基线提供零侵入数据源。
Step 1.2: Polars 外层 Fund 隔离与长缺失切断¶
-
S2 入口 Schema 强断言:读取 Step 1.1 产物后,必须强制比对物理列名(
fund_id,net_value_date,cumulative_net_value)与数据类型(Utf8,Date,Float64),不一致立即抛出AssertionError阻断流水线。 -
长缺失切断判定:按
fund_id分组,计算相邻两行真实数据周日期差值。当差值(周)>MAX_GAP_WEEKS + 1(即 > 5)时,使用累加逻辑(cumsum)为每个连续物理块打上递增的segment_id。单基金内segment_id从 0 起始,各基金独立计算。严禁丢弃任何已有数据行。 -
结构化性能遥测挂载:主执行函数必须强制挂载
core.telemetry.profile_step("S1.2")装饰器,执行完毕后向stderr输出含duration_sec与peak_rss_mb的标准 JSON。 -
统筹编排器物理接入:本步骤由
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 步强制顺序¶
-
微型连续日历生成:按
["fund_id", "segment_id"]提取起止日期,利用pl.int_ranges与pl.duration纯数学加法生成无缺失周序列,与原始表执行 Left Join,并打上is_original标志。 -
首行越界防御:在 Segment 闭包内执行
forward_fill()后,提取每个 Segment 的绝对首行。若cumulative_net_value为 NULL,立即抛出core.exceptions.SegmentFirstRowNullError阻断流水线。 -
血缘切断重算收益率:基于纯净 ffill 序列,使用
shift(1)向量化重算全量weekly_return。严禁复用任何上游历史收益率数据。 -
占位行收益率置 NULL:执行
when(is_original).then(weekly_return).otherwise(None)掩码,将is_original=False的行收益率强制置为 NULL,彻底切断因 ffill 导致的0.0幻影收益向下游rolling_std的传播。 -
结构化性能遥测挂载:主执行函数必须强制挂载
core.telemetry.profile_step("S1.3")装饰器。 - 统筹编排器物理接入:本步骤由统筹器显式调用,接收 Step 1.2 的内存 DataFrame,返回 6 列内存 DataFrame 供下游 Step 1.4 消费。
Step 1.4: 按年切片落盘与 Schema 红线门禁¶
- 年份集合提取:从 Step 1.3 产出的全量 DataFrame 中提取不重复年份并排序。
-
逐年切片落盘:循环内按年份
filter,调用PhysicalStorageContract.sink_parquet_strict()落盘为s2_YYYY.parquet。内部自动注入物理排序(["net_value_date", "fund_id"])与压缩参数(zstd/row_group_size=100000)。 -
循环内就地释放:每次落盘后
del yearly_df,降低进程 RSS 基线,为下游 S3 按年流式加载提供物理前提。 -
Schema 红线门禁:每次落盘后调用
PhysicalStorageContract.validate_schema_symmetry(),读取 Parquet footer metadata(零数据行加载),与S2_OUTPUT_COLUMNS做对称差校验,差集非空直接抛出SchemaDriftError阻断循环。 -
结构化性能遥测挂载:主执行函数必须强制挂载
core.telemetry.profile_step("S1.4")装饰器。 - 统筹编排器物理接入:由统筹器显式调用,接收 Step 1.3 的内存 DataFrame 与输出目录路径,返回已落盘文件路径列表。
- 设计备忘:Step 1.4 的核心价值不在于拦截 S2 阶段自身的内存波峰(S2 全量仅 ~270 万行 × 6 列,Polars 稀疏编码下约 30-50 MB),而在于将 S2 产出物理切分为按年分区文件,防止下游 S3 特征工程(38 维浮点列 × 270 万行)一次性全量加载导致 RSS 飙升至 5-8 GB 的真实雪崩。
第 2 阶段:S3 - 特征工程宏观黑盒¶
全链路算力雪崩最高危区。严禁将本阶段拆分为多个独立 Step 文件串联,必须走“宏观内聚黑盒 -> 按年安全池 -> 纯 NumPy 微线程 -> 物理哨兵封堵”极限拓扑。
Step 2.0: S3 特征工程统一计算与落盘闭环¶
-
宏观黑盒封装:建立唯一的对外入口
src/pipelines/s3/step_s3_feature_engineering.py。统筹器仅需传入输入目录与输出目录,内部封装完整的“按年循环 -> 打散 -> 计算 -> 合并 -> 落盘”生命周期,严禁在统筹器中传递巨大的中间态内存字典。 -
按年流式安全沙箱:外层循环严格按年份读取 S2 产出的单年份物理文件(
s2_YYYY.parquet),单次内存驻留强制锁定在 < 3GB 绝对安全红线内。 -
局部豁免桥接打散:在单年份安全池内,触发
partition_by(["fund_id", "segment_id"], as_dict=True)局部豁免权,将数据零拷贝打散为字典视图,作为向 C 层传递连续内存的桥接器。 -
纯 NumPy 微线程分发:实例化
ThreadPoolExecutor(核数依环境而定)分发字典值。GIL 防线:内层函数第一步必须to_numpy().flatten()抽出纯连续内存数组,后续计算强制约束为仅调用np.dot、np.sum等基础 C 算子,利用底层短暂释放 GIL 实现真并行。 -
纯算子下沉与分母隔离:所有 A/B/C/D/E 类特征计算逻辑作为纯函数下沉至
src/core/math/供单测直接锤炼。所有带窗口期的特征计算,必须绝对遵从“分母隔离强制防线”:遇 NULL 必须直接输出 NULL,严禁残缺窗口计算。 -
物理哨兵生成:在 38 维特征合并后,立即执行全横向非空判定生成
is_feature_complete列。 -
红线门禁与落盘:复用全局
PhysicalStorageContract.sink_parquet_strict()拓扑产出s3_YYYY.parquet。落盘后立即调用validate_schema_symmetry()比对文件 Schema 与 S3 硬编码契约(3 列物理主键 + 38 维特征列 + 1 列哨兵)的对称差,差集非空直接阻断循环并报错。 -
内存就地释放:年份计算完毕落盘后,立即
del释放年份内存池,进入下一年迭代。 -
结构化性能遥测挂载:主执行函数必须强制挂载
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
-
A类 均值回归 (5):
-
1 列物理哨兵:
is_feature_complete(Boolean)
第 3 阶段:S4 - 标签生成¶
业务逻辑核心,防“毒特征污染”与“索引错位”。
Step 3.1: S4 绝对隔离加载¶
-
scan_parquet扫描 S2 产物,立即.select(["fund_id", "net_value_date", "segment_id", "cumulative_net_value"])。严禁拉入weekly_return。
Step 3.2: DCA 闭式向量化公式¶
- 断言
nv_array纯净且大于 0。 - 倒数累加范式:
inv_nv -> cum_inv -> weeks -> returns,全程无循环。
Step 3.3: 止盈边界与 0-based 转 1-based 拦截¶
-
np.where(returns >= 0.20)[0]命中检索。 - 索引偏移修正:
label = int(hit_indices[0] + 1) if len > 0 else 0。 - 绝对边界断言:
assert 0 <= label <= 150。
Step 3.4: 按年落盘¶
- 复用全局工具函数产出
s4_YYYY.parquet。 -
红线 2 拦截落地:落盘后立即比对文件 Schema 与 S4 硬编码契约
{"fund_id", "net_value_date", "segment_id", "label"}的对称差,差集非空直接阻断循环并报错。
第 4 阶段:S5 - 标准化与 Spearman 降维¶
融合“截面数学推导”与“内存防爆累加”。
Step 4.1: 截面聚合与双形态参数物理时序隔离¶
-
纯横截面聚合:按
net_value_date分组计算 38 维的mean, std, median, mad,推导 MAD 上下界,生成is_production_ready哨兵。 -
路径 A (训练宽表):宽表上
.shift(1)实现 T-1 对齐,按日期排序落盘。 -
路径 B (推理长表):Unpivot 为长表,动态生成自然日历
forward_fill至每日,按["feature_name", "date"]绝对物理排序落盘。
Step 4.2: 标准化 Apply 与 IEEE 754 拦截¶
- 实体表
.collect()后与几百行宽表参数 Join(零内存膨胀广播)。 - 遍历 38 列构建
when(std==0).then(None).otherwise(...)表达式,封杀除零产生Inf。
Step 4.3: 训练 Join 三键死锁与毒样本截断¶
- 三键强锁 Join:
["fund_id", "net_value_date", "segment_id"]。 - 哨兵覆写:
is_feature_complete == False的行,标签强制覆写为0。
Step 4.4: Spearman 秩相关内存防爆防线¶
- 按年流式加载标准化特征。
- 截面内
rankdata(method="average")锁定平均秩,算秩后立即减均值。 - Pearson 展开项
np.dot联合累加sum_xy,利用对称性减半算力。 - 跨年累加后执行
(matrix + matrix.T) / 2.0强制对称化截断。 - 剔除
>= 0.90的冗余特征,输出final_feature_list。
第 5 阶段:S5.3 - 训练流水线¶
将降维后的纯净特征转化为 LightGBM LambdaRank 可消费的形态并进行闭环重训。
Step 5.1: LambdaRank 样本平衡重构¶
-
Group 划分:全量表按
net_value_date排序,每个日期即一个 Query Group。 -
差异化等距抽样:正样本池(
label>0)用np.linspace抽最多 50 条;负样本池抽最多 15 条。 -
相关性得分映射:
label=0映射为0,label=1~150映射为151 - label(第 1 周得 150 分)。 -
Group 边界构建:记录每个截面抽样数,构建一维数组传给 LightGBM 的
group参数。
Step 5.2: Time-Series 5 折评估防线¶
- 时间切分:按 20% 步长划分 5 个验证区块,使用累积扩展窗口。
-
NDCG@10 绝对核心:训练传入
group_boundaries,评估强制锁定NDCG@10指标,直接对齐 Top 10 榜单业务 KPI。 - 防穿透锁定:5 折期间严禁再次剔除特征或修改任何超参,仅观察指标分布。
Step 5.3: 闭环全量重训与产物契约落盘¶
- 使用全量抽样数据 +
final_feature_list进行最终无截断训练。 -
产物强制打包:将
model.txt与两张标准化参数表打入同一个带时间戳的发布目录(如artifacts/v_20231024/)。 - MD5 强校验:生成模型文件的 MD5 摘要,作为后续 S6 加载前的防篡改门禁。
第 6 阶段:S6 - 推理服务¶
纯粹的“无状态数学计算引擎”,彻底剥离所有 I/O 职责。
Step 6.1: FastAPI 骨架与进程级全局缓存防抖¶
-
生命周期加载:在
lifespan启动钩子中,读取发布目录产物,加载_PARAMS_LONG_CACHE(长表)、_AVAIL_DATES_CACHE(日期数组)、_MODEL(Booster)。 -
MD5 门禁校验:加载模型前比对 MD5,不一致直接
os._exit(1)拒绝启动。
Step 6.2: 两步走时间锚定防线¶
-
第一步锚定:接收单维
request_date,与_AVAIL_DATES_CACHE执行极速asof_join寻找valid_date。无效则短路返回空结构体。 -
第二步拉取:基于
final_feature_list构建左表,按["feature_name", "valid_date"]排序,与_PARAMS_LONG_CACHE执行双键asof_join,后置过滤is_production_ready。
Step 6.3: 2D 矩阵级向量化批处理内核¶
-
统一抽象:无论单基金(N=1)还是批量(N=1000),统一接收
(N, 38)原始特征 2D 矩阵。 -
复用 S5 内核:调用
core.standardization_math.apply_zscore,利用 NumPy 广播一次性完成 N 行的 MAD 裁剪与 Z-Score 转换。 -
矩阵预测:
_MODEL.predict(matrix),利用 C++ 底层极速输出得分。
Step 6.4: 接口物理隔离与归因红线¶
-
单基金接口:允许调用
predict_contrib返回 SHAP 归因明细。 -
批量接口:架构级禁止引入
predict_contrib代码,归因字段强制硬编码为空字符串""。
第 7 阶段:端到端集成与打通¶
这是将离线数学模型转化为线上业务价值的最终闭环,解决“产物交接、预热、回测对接、极限压测”四大工程鸿沟。
Step 7.1: 产物自动化交接与灰度发布契约¶
-
CI/CD 打包流:S5 训练流水线成功后,CI 自动将
artifacts/v_xxx/打包为 Docker Image 或推送到 OSS/Minio。 - S6 滚动更新:K8s/ Docker Compose 执行滚动更新。新 Pod 启动时拉取最新产物进行 MD5 校验与预热,校验失败则 Pod 启动失败,旧 Pod 继续服役,实现无损回滚。
Step 7.2: S6 冷启动掩盖与 Health Check 强断言¶
-
就绪探针:在 FastAPI 中实现
/health接口。不仅返回 200,内部必须验证_MODEL is not None且_PARAMS_LONG_CACHE.shape[0] > 0。未加载完毕前,K8s 绝不将流量打入该 Pod。 - 预热请求:Pod 启动后,后台线程自动发起一次包含 10 只基金的矩阵预测请求,强制触发 LightGBM 内部树的 JIT 缓存预热,掩盖首次真实请求的微秒级延迟毛刺。
Step 7.3: 回测引擎标准化数据对接¶
- 定义回测接口契约:S6 批量接口不直接服务 C 端,而是作为 BFF 层被回测引擎(如 Backtrader/Qlib)调用。
-
历史截面回放:回测引擎按历史日期循环,将历史 T 日全市场基金原始特征组装为
(N, 38)矩阵,调用 S6 获取 T 日得分排名。 - 严格按照 NDCG@10 验证:回测引擎拿到 Top 10 池后,向后看 150 周计算真实定投收益率。对比回测出的“Top 10 平均达标率”与 S5 验证集的“NDCG@10 曲线”,两者必须呈现强正相关,否则证明存在数据穿越 Bug。
Step 7.4: 全链路极限压测与 cgroup 防雪崩验证¶
- 流量炮台构建:使用 Locust 编写压测脚本,模拟 50 并发线程,每线程持续发送 N=1000 的批量推荐请求。
-
P99 延迟红线:监控 S6 服务的 P99 响应时间,必须死锁在
< 50ms以内。一旦超限,说明触犯了矩阵广播外的隐性锁或内存交换。 -
内存防雪崩观察:在压测期间,通过
docker stats严格监视 S6 容器的内存曲线。因为 S6 是纯矩阵运算无中间态膨胀,其内存曲线必须是一条绝对水平的直线。任何呈阶梯状上升的现象,都证明存在隐藏的内存泄漏或未触发全局缓存的灾难逻辑,必须阻断上线。
由 Huarui Lin 更新于 5 个月 之前 · 13 修订