从单个变换到可审计系统
上一章回答了『特征如何变换』,本章进一步回答『上百个变换如何在每次训练、验证和上线时保持同一语义』。规模化不只是增加 CPU 核数;它首先是关于列语义、拟合边界、数据布局、稀疏性和资源预算的工程约束。
本章先修是第 15 章的标准化、缩尾、类别编码与数据泄漏,以及基本二元分类和交叉验证。学完后,读者应能够:
解释 Pipeline 与 ColumnTransformer 如何绑定变换顺序、列选择和拟合状态;
用本地真实估值数据构建数值列与类别列的端到端流水线;
用扩展时间窗口评估模型,按标签可得日 purge 每个拟合折,并保证所有拟合型转换仅从折内训练样本估计;
分析稠密与稀疏布局、复制成本、调度开销、内存峰值和并行过度;
设计可重现的并行交叉验证与参数搜索,并为资源边界设置防护。
本章的正式目标通过下列活动和核心答案逐项评价。
表 表 16.1 汇总本节的计算或审计结果,解释时应遵循正文给出的口径与限制。
真实估值数据与预测时点
案例目标是在季末已观测的估值与公司属性上,识别严格下一季度市值变化率位于同期截面前五分位的长三角公司。列表 16.1 建立本地数据接口,列表 16.2 显式保留 feature_date、target_date 与 label_available_date。由于市值是市场日数据,标签可得日取目标季度末;若目标来自财报,应改用实际披露日。这里的逻辑回归只是检验端到端接口的可执行下游,不构成投资策略结论或统计显著性结论。
易混淆概念辨析:市值变化率的公司行动边界
相邻季度总市值之比同时混合价格变化与总股本变化;送转、增发、回购等公司行动都可能改变分母或分子,现金分红也未进入该比率。由于当前估值表没有逐期股份数量与公司行动调整链,本章统一使用『下一季度市值变化率/截面前五分位』,不把目标称为股票收益、超额收益或显著跑赢。
表 16.2 披露筛选前的预测跨度,并断言进入目标面板的行全部严格相隔三个月且市值变化率可计算。特征列包括四个估值比率、对数市值,以及行业与省份。2025 年第一个有效季度末定义最终测试信息时点;候选历史特征即使早于该日,只要标签在该日或之后才形成,也必须从最终训练中 purge。表 16.3 显示这一时间留出的实际时点、purge 数和非法标签数。
prediction_horizon_audit['进入目标面板数'] = 0 # 明确其他自然月跨度不进入同一分类任务
prediction_horizon_audit.loc[3, '进入目标面板数'] = len(model_panel) # 记录三个月且市值变化率可计算的样本
prediction_horizon_audit['排除数'] = prediction_horizon_audit['筛选前样本数'].sub(prediction_horizon_audit['进入目标面板数']) # 量化跨度或目标缺失造成的剔除规模
assert model_panel['horizon_months'].eq(3).all() # 防止下一观测期再次替代固定下一季度
prediction_horizon_audit # 展示筛选前后各预测跨度的数量
numeric_feature_names = ['pe_ratio_ttm', 'pb_ratio_lf', 'ps_ratio_ttm', 'dividend_yield_ttm', 'log_market_cap'] # 定义需要插补、缩尾和标准化的列
categorical_feature_names = ['industry', 'province'] # 定义需要众数插补和 One-Hot 编码的列
model_feature_names = numeric_feature_names + categorical_feature_names # 固定模型接口的列名与顺序
test_information_start = model_panel.loc[model_panel['feature_date'] >= pd.Timestamp('2025-01-01'), 'feature_date'].min() # 以首个2025年特征截面作为最终决策时点
assert pd.notna(test_information_start) # 保证本地数据提供可用时间留出期
is_historical_feature = model_panel['feature_date'] < test_information_start # 标识测试时点前的候选历史特征
is_label_known_before_test = model_panel['label_available_date'] < test_information_start # 标识测试决策前严格可得的历史结果
is_training_row = is_historical_feature & is_label_known_before_test # 按标签可得日 purge 最终拟合集
is_test_row = model_panel['feature_date'] >= test_information_start # 保留测试信息时点及之后的样本
purged_holdout_boundary_count = int((is_historical_feature & ~is_label_known_before_test).sum()) # 统计边界处尚无结果的候选训练行
training_panel = model_panel.loc[is_training_row].copy() # 保存全部时点列以供后续折级审计
test_panel = model_panel.loc[is_test_row].copy() # 保存未来特征和事后评价结果
holdout_illegal_label_count = int(training_panel['label_available_date'].ge(test_information_start).sum()) # 计算最终训练中的非法标签数
assert holdout_illegal_label_count == 0 # 阻止留出期结果进入 Pipeline、CV 和调参
training_features = training_panel.loc[:, model_feature_names].copy() # 从安全训练样本提取固定特征接口
test_features = test_panel.loc[:, model_feature_names].copy() # 从留出样本提取相同特征接口
training_target = training_panel['is_top_quintile'].copy() # 取得安全训练样本的二元目标
test_target = test_panel['is_top_quintile'].copy() # 取得只用于事后评分的留出目标
split_audit = pd.DataFrame({'特征起始日': [training_panel['feature_date'].min(), test_panel['feature_date'].min()], '特征结束日': [training_panel['feature_date'].max(), test_panel['feature_date'].max()], '最大标签可得日': [training_panel['label_available_date'].max(), test_panel['label_available_date'].max()], '样本数': [len(training_panel), len(test_panel)], '边界 purge 数': [purged_holdout_boundary_count, 0], '非法训练标签数': [holdout_illegal_label_count, 0], '正类比例': [training_target.mean(), test_target.mean()]}, index=['最终训练集', '2025 时间留出集']) # 检查时间、规模和标签可得性
split_audit # 将最终训练和留出边界显示为表格
Pipeline 的状态机制
对两步流水线 \(f_{β}∘T_{θ}\),fit(X_train, y_train) 依次估计变换参数 \(θ\) 和模型参数 \(β\);predict(X_test) 只计算 \(f_{β}(T_{θ}(X_{test}))\)。交叉验证每一折都会 clone 出一个未拟合估计器,从而使该折的验证数据无法进入中位数、分位数、均值或词表的估计。
易混淆概念辨析:Pipeline 只保护放进它内部的操作
如果研究者在调用 Pipeline 之前已用全样本填补缺失、筛选特征或做目标编码,Pipeline 无法追溯并撤销这些泄漏。标签日期也在 Pipeline 之外:必须先按 label_available_date 构造安全索引,再让 Pipeline 只接触这些索引对应的特征和目标。安全单元必须同时包含正确的样本信息集与每个会从数据学习参数的步骤。
可拟合的缩尾变换器
下面的 QuantileClipper 把训练分位数保存为带下划线的拟合属性。这一设计使其可被 sklearn 克隆、交叉验证和参数搜索,完整实现见 列表 16.3。
时间交叉验证
普通 K 折将未来日期混入训练折,不符合投资决策的信息顺序。即使候选拟合行满足 feature_date < v_j,其三个月后标签也可能在验证日 \(v_j\) 当天或之后才形成。因而第 \(j\) 折先建立历史特征候选集,再 purge 所有 label_available_date >= v_j 的行,并用单一日期 \(v_j\) 验证。验证标签仅在结果形成后用于评分,不参与该折拟合。每折的 Pipeline 都是新克隆,所以中位数、缩尾边界和 One-Hot 词表也只由安全拟合端估计;实际折边界见 表 16.5。
training_dates = training_panel['feature_date'].reset_index(drop=True) # 保留与 training_features 相同的特征日顺序
training_label_available_dates = training_panel['label_available_date'].reset_index(drop=True) # 保留每条训练目标的可得日
unique_training_dates = np.sort(training_dates.unique()) # 提取可作为完整截面的历史日期
validation_dates = unique_training_dates[-3:] # 选取训练期最后三个截面评估稳定性
date_cv_splits = [] # 收集 sklearn 可直接使用的位置索引对
fold_rows = [] # 收集人可读的时间与样本审计记录
for validation_information_time in validation_dates: # 依时间顺序构造扩展窗口
is_fit_candidate = training_dates < validation_information_time # 先建立特征日严格早于验证日的候选集
is_fit_label_known = training_label_available_dates < validation_information_time # 再要求结果在验证决策前已经可得
fit_positions = np.flatnonzero((is_fit_candidate & is_fit_label_known).to_numpy()) # purge 候选集中尚未形成的标签
validation_positions = np.flatnonzero(training_dates.to_numpy() == validation_information_time) # 保持整个当期截面作为验证折
purged_label_count = int((is_fit_candidate & ~is_fit_label_known).sum()) # 量化只审计特征日会误纳入的标签
illegal_label_count = int(training_label_available_dates.iloc[fit_positions].ge(validation_information_time).sum()) # 计算最终拟合端非法标签数
assert illegal_label_count == 0 and len(fit_positions) > 0 # 要求每折有历史样本且非法标签严格为零
date_cv_splits.append((fit_positions, validation_positions)) # 向 sklearn 交付明确的无未来分割
fold_rows.append({'验证信息日': validation_information_time, '候选拟合数': int(is_fit_candidate.sum()), '标签 purge 数': purged_label_count, '最终拟合数': len(fit_positions), '最大标签可得日': training_label_available_dates.iloc[fit_positions].max(), '非法标签数': illegal_label_count, '验证样本数': len(validation_positions)}) # 记录折级时点与样本变化
time_fold_audit = pd.DataFrame(fold_rows) # 汇总每折标签可得性证据
assert time_fold_audit['非法标签数'].eq(0).all() # 防止后续修改重新引入未来结果
time_fold_audit # 展示每折 purge 数与非法标签零计数
Joblib 并行与开销模型
设一个任务的有效计算时间为 \(T_{work}\),序列化、进程调度、数据复制和结果合并的开销为 \(T_{overhead}\)。在 \(p\) 个工作进程下,式 16.1 给出一个粗略上界:
\[
T_p\approx \frac{T_{work}}{p}+T_{overhead}(p).
\tag{16.1}\]
当数据小、折数少或模型很快时,\(T_{overhead}\) 可能大于节省的计算,两核会比单核更慢。因此 表 16.6 的本地计时是实验结果,不是『并行必然加速』的保证。
from sklearn.model_selection import cross_validate # 在每折内克隆并拟合完整流水线
from time import perf_counter # 用单调时钟测量墙上时间
scoring_metrics = {'roc_auc': 'roc_auc', 'average_precision': 'average_precision'} # 同时评估排序质量与罕见正类检出
serial_start = perf_counter() # 记录单核评估的起点
serial_results = cross_validate(feature_pipeline, training_features, training_target, cv=date_cv_splits, scoring=scoring_metrics, n_jobs=1, return_train_score=False) # 以单核建立可比基准
serial_seconds = perf_counter() - serial_start # 计算单核完整墙上时间
parallel_start = perf_counter() # 记录两核评估的起点
parallel_results = cross_validate(feature_pipeline, training_features, training_target, cv=date_cv_splits, scoring=scoring_metrics, n_jobs=2, pre_dispatch=2, return_train_score=False) # 限制预分派任务以控制内存峰值
parallel_seconds = perf_counter() - parallel_start # 计算包含调度开销的并行时间
parallel_audit = pd.DataFrame({'模式': ['单核', '两核'], '墙上秒数': [serial_seconds, parallel_seconds], '平均 ROC-AUC': [serial_results['test_roc_auc'].mean(), parallel_results['test_roc_auc'].mean()], '平均 PR-AUC': [serial_results['test_average_precision'].mean(), parallel_results['test_average_precision'].mean()]}) # 对照时间与数值结果是否一致
parallel_audit.round(4) # 报告当前机器和数据切片下的实测结果
并行参数搜索
参数名通过双下划线沿管线层级定位。例如 classifier__C 表示分类器步骤的逆正则化强度。搜索必须使用同一组经标签可得日 purge 的时间分割,否则并行只是更快地得到有偏结果。候选结果和最终 2025 年评估分别见 表 16.7 与 表 16.8。
from sklearn.model_selection import GridSearchCV # 组合管线参数与扩展窗口评估
parameter_grid = {'classifier__C': [0.1, 1.0], 'classifier__class_weight': [None, 'balanced']} # 使用小型网格比较正则化与类权重
grid_search = GridSearchCV(feature_pipeline, parameter_grid, scoring=scoring_metrics, refit='average_precision', cv=date_cv_splits, n_jobs=2, pre_dispatch=2, return_train_score=False) # 以罕见正类的 PR-AUC 选择并重拟合
assert holdout_illegal_label_count == 0 and time_fold_audit['非法标签数'].eq(0).all() # 在调参前锁定最终训练和全部折的标签边界
grid_search.fit(training_features, training_target) # 在安全历史训练集和安全时间折上执行全部参数比较
search_results = pd.DataFrame(grid_search.cv_results_) # 将每个候选组合的折级结果转为可审计表
search_summary = search_results.loc[:, ['param_classifier__C', 'param_classifier__class_weight', 'mean_test_roc_auc', 'mean_test_average_precision', 'mean_fit_time']].sort_values('mean_test_average_precision', ascending=False) # 聚焦性能与拟合成本的取舍
search_summary.round(4) # 显示全部候选而不只是最优结果
from sklearn.metrics import average_precision_score, roc_auc_score # 对罕见正类同时评估 PR 与 ROC 排序
holdout_probabilities = grid_search.predict_proba(test_features)[:, 1] # 用首个测试时点前已知标签选定的完整流水线预测 2025 年
holdout_results = pd.DataFrame({'指标': ['ROC-AUC', 'PR-AUC', '正类基准比例'], '数值': [roc_auc_score(test_target, holdout_probabilities), average_precision_score(test_target, holdout_probabilities), test_target.mean()]}) # 将 PR-AUC 与无技能基准对照
holdout_results.round(4) # 报告未参与调参的真正时间样本外结果
如果 PR-AUC 仅接近正类基准比例,正确结论是当前特征与线性决策边界的样本外信号弱,而不是继续加大搜索网格直到出现满意数字。
数据布局、复杂度与资源边界
稠密、稀疏与列式布局
若 One-Hot 矩阵有 \(n\) 行、\(p\) 列和 \(m\) 个非零元素,稠密 float64 存储约需 \(8np\) 字节;CSR 稀疏矩阵主要存储 \(m\) 个值、\(m\) 个列索引和 \(n+1\) 个行指针,空间近似与 \(m\) 而不是 \(np\) 成正比。高基数类别默认保持稀疏是关键资源约束。
HDF5 fixed 布局适合一次性整表读取,却不支持谓词下推。大规模生产数据更适合按日期或市场分区的 Parquet/Arrow 列式布局,使数据层能只读取本次特征所需的分区和列。Pipeline 解决变换状态,不代替数据层的投影和谓词下推。
并行资源防护
n_jobs=-1 会尽可能占用所有核,并不代表资源最优;共享服务器上应为其他作业保留容量。
外层 Joblib 进程与内层 BLAS/OpenMP 线程同时展开会产生过度并行;应将其中一层的线程数限制为 1。
进程后端需要序列化估计器和数据;大数组可以通过 Joblib memmap 减少复制,但 Pandas 对象列常仍有明显序列化成本。
pre_dispatch 控制预先创建的任务数。内存受限时,不应让尚未执行的任务同时保留多份数据。
交叉验证的峰值内存大致取决于同时拟合数与单折峰值的乘积。在小样本上测量单折峰值,再增加并发度,比直接使用全部核更稳妥。
常见误区
对全表调用 fit_transform 后再交叉验证:这会把每个验证折的分布信息泄漏给拟合折。
只比较特征日期:即使 fit_feature_date < validation_date,拟合标签仍可能在验证日之后才形成。
抽样后丢失时间顺序:随机 K 折可能在同一折中用未来预测过去。
将稀疏矩阵转为稠密矩阵:一次 .toarray() 可能使内存需求从 \(O(m)\) 跃升到 \(O(np)\)。
认为并行必然线性加速:调度、复制、内存带宽和串行比例都会限制加速。
将搜索最佳折内得分当作最终表现:参数选择已使用了交叉验证信息,仍需未触碰的时间留出集。
本章小结
Pipeline 将变换状态和模型状态封装成一个可克隆对象,ColumnTransformer 则将列语义显式编码为处理图。但它们不会自动理解标签何时形成:必须先按 label_available_date purge 最终训练和每个时间折,再在安全索引上拟合全部变换。并行是由工作量、调度开销、复制和内存预算共同决定的工程选择,而不是一个无条件的加速开关。
分层练习
基础题
练习 16.1
写出从完整管线中访问数值插补中位数、One-Hot 类别和分类器系数的属性路径,并说明未拟合管线为何没有这些带下划线的属性。
完整解答
拟合状态的实际路径见 表 16.9。
best_pipeline = grid_search.best_estimator_ # 取出已在全训练期重拟合的最佳管线
best_preprocessor = best_pipeline.named_steps['preprocessor'] # 定位已拟合的列变换器
numeric_medians = best_preprocessor.named_transformers_['numeric'].named_steps['imputer'].statistics_ # 读取仅由训练期估计的中位数
encoded_categories = best_preprocessor.named_transformers_['categorical'].named_steps['encoder'].categories_ # 读取训练期出现的类别词表
classifier_coefficients = best_pipeline.named_steps['classifier'].coef_ # 读取在变换特征上拟合的线性系数
state_summary = pd.DataFrame({'状态': ['数值中位数个数', '类别列数', '分类器系数个数'], '值': [len(numeric_medians), len(encoded_categories), classifier_coefficients.shape[1]]}) # 检查三类拟合状态的维度
state_summary # 展示流水线各层的可审计状态
带下划线的 sklearn 属性是 fit 产生的学习结果。构造器只记录超参数,在看到训练数据前不应已经拥有中位数、类别表或系数。
进阶题
练习 16.2
对每个时间折同时验证 max(fit_feature_date) < validation_information_time 与 max(fit_label_available_date) < validation_information_time,报告被 purge 的候选标签数、非法标签数,并检查验证折是否包含完整的单一截面。解释为何只检查特征日期会错误放行尚未形成的标签。
完整解答
完整断言见 列表 16.5。
exercise_fold_rows = [] # 收集练习要求的折级证据
for fold_number, (fit_positions, validation_positions) in enumerate(date_cv_splits, start=1): # 逐折审计交给 sklearn 的实际索引
fit_feature_dates = training_dates.iloc[fit_positions] # 恢复拟合端特征日
fit_label_dates = training_label_available_dates.iloc[fit_positions] # 恢复拟合端标签可得日
fold_validation_dates = training_dates.iloc[validation_positions] # 恢复当前折验证信息日
validation_information_time = fold_validation_dates.min() # 取得该完整截面的单一决策时点
invalid_feature_count = int(fit_feature_dates.ge(validation_information_time).sum()) # 统计特征日期违规行
invalid_label_count = int(fit_label_dates.ge(validation_information_time).sum()) # 直接统计当时尚不可得的拟合标签
single_date_count = int(fold_validation_dates.nunique()) # 检查验证折是否保持完整单日截面
purged_label_count = int(time_fold_audit.loc[time_fold_audit['验证信息日'] == validation_information_time, '标签 purge 数'].iloc[0]) # 读取构造折时实际剔除规模
exercise_fold_rows.append({'折': fold_number, '验证信息日': validation_information_time, '标签 purge 数': purged_label_count, '非法特征数': invalid_feature_count, '非法标签数': invalid_label_count, '验证日期数': single_date_count}) # 保存可评分审计结果
exercise_time_audit = pd.DataFrame(exercise_fold_rows) # 形成练习答案审计表
assert exercise_time_audit[['非法特征数', '非法标签数']].eq(0).all().all() # 要求拟合特征和标签均严格早于验证时点
assert exercise_time_audit['验证日期数'].eq(1).all() # 要求每折验证集是完整单一截面
exercise_time_audit # 展示各折 purge 后非法标签数为零
只检查特征日期会把『早先形成的特征、尚未形成的三个月后结果』误当作可训练样本。练习表中 标签 purge 数 可以大于 0,因为它记录修正前的候选污染;进入实际拟合折后的 非法标签数 必须恒为 0。
真实数据应用题
练习 16.3
计算测试特征稀疏矩阵的密度,并比较 CSR 核心数组与理论稠密 float64 矩阵的字节数。不得调用 .toarray()。
完整解答
存储估算的可执行答案见 表 16.10。
assert sparse.issparse(transformed_test_features) # 保证未在审计前意外稠密化特征矩阵
matrix_rows, matrix_columns = transformed_test_features.shape # 读取特征矩阵的逻辑维度
matrix_density = transformed_test_features.nnz / (matrix_rows * matrix_columns) # 计算非零元素占全矩阵的比例
csr_bytes = transformed_test_features.data.nbytes + transformed_test_features.indices.nbytes + transformed_test_features.indptr.nbytes # 加总 CSR 值、列索引和行指针存储
dense_bytes = matrix_rows * matrix_columns * np.dtype('float64').itemsize # 在不创建稠密矩阵的前提下计算理论字节数
memory_comparison = pd.DataFrame({'布局': ['CSR 实际核心数组', '稠密 float64 理论值'], '字节数': [csr_bytes, dense_bytes], '矩阵密度': [matrix_density, 1.0]}) # 对照两种布局的资源代价
memory_comparison.round(4) # 报告实际稀疏性与存储差异
综合挑战题
练习 16.4
现有 24 个参数组合和 5 个时间折,单折拟合峰值内存实测为 1.8 GB,可用内存为 12 GB,并为系统保留 3 GB。给出保守的 n_jobs 和 pre_dispatch 设置,计算总拟合数,并说明为何参数组合数不等于峰值并发数。
完整解答
可用于拟合的内存为 \(12-3=9\) GB。根据单折 1.8 GB,纯算术上限为 \(⌊9/1.8⌋=5\) 个同时拟合;考虑 Python 进程、数据复制和短时峰值,保守选择 n_jobs=4, pre_dispatch=4。总拟合数为 \(24×5=120\),另加最优管线在全训练集上的一次重拟合。这 120 个任务排队执行,参数组合数只决定任务总数;n_jobs 和 pre_dispatch 才限制一个时点同时存活或预分派的任务数。上线前还应在代表性数据分区上监测 RSS 峰值,而不只依赖算术估计。