Course 1:ML 流水线构建
课程简介
特征工程、数据流水线、模型训练工作流。
🎬 本课程视频:MLOps Production — 机器学习工程生产实践
一、特征工程
1.1 特征工程的核心目标
特征工程是将原始数据转换为模型可以高效学习的特征的过程。它是 ML 项目中最重要也最耗时的环节之一。
核心目标:
1. 提升预测能力:创建能捕捉数据中重要模式的特征
2. 保证可解释性:特征应该具有实际含义
3. 减少噪声:过滤掉对预测无帮助或有害的信息
1.2 数值型特征处理
缩放:
- 归一化(Min-Max):将特征缩放到 [0, 1] 范围
- 标准化(Z-Score):减去均值除以标准差
- 鲁棒缩放:减去中位数除以 IQR
截断(Clipping):对异常值进行上下界截断
非线性变换:
- 对数变换:$\log(x)$ 处理长尾分布
- 平方根:$\sqrt{x}$
- Box-Cox 变换:自动寻找最优幂变换
离散化:将连续值分段
- 等宽分桶:将范围分成等宽区间
- 等频分桶:每桶包含相同数量的样本
- 自定义分桶:基于业务知识的分段
1.3 类别型特征处理
独热编码(One-Hot Encoding):将 K 个类别转为 K 个二元特征
标签编码:将 K 个类别映射为 1, 2, ..., K(适合树模型)
目标编码:用类别对应的目标均值替换原始类别值(注意防止过拟合)
频次编码:用类别在数据中出现的频率替换类别
# 独热编码
from sklearn.preprocessing import OneHotEncoder
encoder = OneHotEncoder(sparse=False, handle_unknown='ignore')
X_encoded = encoder.fit_transform(X_categorical)
# 目标编码
from category_encoders import TargetEncoder
encoder = TargetEncoder()
X_encoded = encoder.fit_transform(X, y)
1.4 时序特征处理
滞后特征(Lag Features):使用过去时刻的值
- $x_{t-1}, x_{t-2}, ..., x_{t-k}$
滑动窗口统计:
- 均值、标准差、最大值、最小值、趋势
时间特征:
- 年、月、日、周几、小时
- 节假日标记
- 季节分量
def create_lag_features(df, column, lags=[1, 7, 30]):
for lag in lags:
df[f'{column}_lag_{lag}'] = df[column].shift(lag)
return df
def create_rolling_features(df, column, windows=[7, 30]):
for w in windows:
df[f'{column}_rolling_mean_{w}'] = df[column].rolling(w).mean()
df[f'{column}_rolling_std_{w}'] = df[column].rolling(w).std()
return df
1.5 特征选择
过滤法:基于统计指标筛选
- 方差过滤:方差过小的特征无信息量
- 相关系数:与目标相关性高的特征
- 互信息:捕捉非线性关系
包裹法:基于模型效果选择
- 前向选择:逐步添加最有帮助的特征
- 后向消除:逐步删除最无帮助的特征
- 递归特征消除(RFE)
嵌入法:在模型训练中完成选择
- L1 正则化(Lasso):自动稀疏化
- 树模型特征重要性
1.6 特征存储(Feature Store)
特征存储是 MLOps 的重要基础设施:
- 在线存储:低延迟的特征服务(Redis、DynamoDB)
- 离线存储:批量特征计算(HDFS、S3)
- 特征注册:特征的元数据管理、发现、复用
常用工具:Feast、Tecton、Databricks Feature Store
二、数据管道
2.1 ETL 与 ELT
ETL(Extract, Transform, Load):
1. 从源系统抽取数据
2. 在中间层进行转换
3. 加载到目标存储
ELT(Extract, Load, Transform):
1. 从源系统抽取数据
2. 直接加载到目标存储(如数据湖)
3. 在查询时进行转换
现代大数据架构倾向于 ELT,利用数据湖的计算能力。
2.2 批处理与流处理
批处理(Batch Processing):
- 定期执行(每天/每小时)
- 处理大量历史数据
- 工具:Apache Spark、Airflow、dbt
- 适合:特征批量计算、模型定期训练
流处理(Stream Processing):
- 实时处理到达的数据
- 低延迟(毫秒-秒级)
- 工具:Apache Kafka、Flink、Spark Streaming
- 适合:实时特征、在线预测、监控告警
2.3 管道编排
工作流调度:
- Apache Airflow:DAG 驱动的任务编排
- Prefect:现代化的工作流管理
- Dagster:面向数据资产的工作流
# Airflow DAG 示例
from airflow import DAG
from airflow.operators.python import PythonOperator
with DAG('ml_pipeline', schedule_interval='@daily') as dag:
extract = PythonOperator(task_id='extract', python_callable=extract_data)
transform = PythonOperator(task_id='transform', python_callable=transform_data)
train = PythonOperator(task_id='train', python_callable=train_model)
evaluate = PythonOperator(task_id='evaluate', python_callable=evaluate_model)
extract >> transform >> train >> evaluate
2.4 数据质量监控
在管道中嵌入数据质量检查:
- 完整性检查:字段缺失率是否在允许范围内
- 唯一性检查:主键是否唯一
- 范围检查:数值是否在合理范围内
- 分布检查:数据分布是否显著变化
三、训练工作流
3.1 实验追踪
每次训练实验应该记录:
- 数据集版本和哈希值
- 模型超参数
- 训练/验证指标
- 模型 artifacts
- 环境和依赖版本
工具:MLflow、Weights & Biases、Neptune
import mlflow
with mlflow.start_run():
mlflow.log_param("learning_rate", 0.01)
mlflow.log_param("n_estimators", 100)
mlflow.log_metric("accuracy", 0.95)
mlflow.log_artifact("model.pkl")
3.2 超参数调优
手动搜索:基于经验手动调整
网格搜索:穷举所有参数组合
随机搜索:从参数分布中随机采样
贝叶斯优化:基于历史结果构建概率模型,指导下一组参数选择
from sklearn.model_selection import RandomizedSearchCV
param_dist = {
'n_estimators': [50, 100, 200],
'max_depth': [None, 10, 20, 30],
'learning_rate': [0.01, 0.05, 0.1]
}
search = RandomizedSearchCV(model, param_dist, n_iter=20, cv=5)
search.fit(X_train, y_train)
3.3 模型注册
模型注册表管理模型的完整生命周期:
- 模型版本控制
- 模型元数据(来源、性能、用途)
- 模型状态(开发、预发布、生产、退役)
- 审批流程
四、总结
- 特征工程是将原始数据转换为模型特征的关键步骤
- 特征存储统一管理在线/离线特征
- 数据管道通过编排 ETL/ELT 流程确保数据持续可用
- 实验追踪记录训练过程的完整信息,保证可复现性
- 超参数调优系统性地搜索最优配置
五、特征工程进阶
5.1 特征交叉
特征交叉(Feature Crossing)是将两个或多个特征组合起来创建新特征的方法。它可以帮助线性模型学习非线性关系。
例如,在广告点击率预测中,交叉"用户年龄"和"广告类别"可以创建更有信息量的特征——年轻用户可能对游戏广告感兴趣,而中年用户可能对理财产品感兴趣。
在 scikit-learn 中,PolynomialFeatures 和实现特征交叉:
from sklearn.preprocessing import PolynomialFeatures
# 创建 [x1, x2, x1*x2] 特征
cross = PolynomialFeatures(degree=2, interaction_only=True)
X_cross = cross.fit_transform(X[['age', 'ad_category']])
5.2 自动化特征工程
在大规模 ML 系统中,手动创建所有特征是不现实的。自动化特征工程工具可以系统性地探索特征空间:
- Featuretools:基于"深度特征合成"(Deep Feature Synthesis)的自动化特征工程
- TSFresh:自动提取时序特征
- AutoFeat:通过遗传编程搜索特征变换
5.3 特征重要性分析
理解哪些特征对模型最重要,可以帮助:
- 减少特征数量(节省存储和计算)
- 提供业务洞察(哪些因素驱动预测结果)
- 检测数据质量问题(最重要的特征大量缺失)
树模型直接提供特征重要性。对线性模型,系数绝对值可以作为重要性的度量。对任意模型,可以通过排列重要性(Permutation Importance)或 SHAP 值来分析。
5.4 数据管道监控
数据管道需要监控以下关键指标:
- 数据新鲜度:数据从产生到可用需要多久
- 数据完整性:数据是否有缺失或延迟
- 特征覆盖率:每个特征的非空比例
- 管道失败率:管道运行的失败次数和原因
建立数据管道监控仪表盘,让团队能实时了解数据管道的健康状况。
特征存储(Feature Store)
在成熟的 ML 平台中,特征存储(Feature Store)是一个关键的架构组件。Feature Store 解决了特征管理中的几个核心问题。第一是特征复用——不同模型团队可以共享和复用加工好的特征,避免重复计算。第二是在线-离线一致性——训练时使用的特征工程逻辑与推理时使用的逻辑完全一致,这是 ML 流水线中最容易出问题的环节。
Feature Store 通常包含两个部分:离线 Feature Store(存储大规模历史特征数据,通常基于 Parquet 或 ORC 格式,存放在数据湖中)和在线 Feature Store(存储低延迟的实时特征数据,通常基于 Redis 或 Cassandra 等 KV 存储)。当训练模型时,从离线 Feature Store 读取特征;当模型在线服务时,从在线 Feature Store 读取最新特征。这种双存储架构确保了训练和推理的一致性,同时满足了不同场景的性能需求。
数据管道中的血缘追踪
数据血缘(Data Lineage)追踪是数据管道的另一个重要功能。血缘追踪记录了从原始数据到最终特征的完整转换过程——每个特征是由哪些原始字段、通过什么转换逻辑生成的。当发现某个特征有问题时,血缘追踪可以帮助我们快速定位问题的来源,并评估受影响的模型和下游任务。现代数据平台如 Apache Atlas、Amundsen、DataHub 都提供了自动化的血缘追踪能力。
延伸阅读
- 📺 B 站播放列表:MLOps Production — 机器学习工程生产实践
- 📚 更多学习资源,请访问 deeplearning.ai 官网