配音

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 的重要基础设施:

常用工具: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 模型注册

模型注册表管理模型的完整生命周期:
- 模型版本控制
- 模型元数据(来源、性能、用途)
- 模型状态(开发、预发布、生产、退役)
- 审批流程

四、总结

  1. 特征工程是将原始数据转换为模型特征的关键步骤
  2. 特征存储统一管理在线/离线特征
  3. 数据管道通过编排 ETL/ELT 流程确保数据持续可用
  4. 实验追踪记录训练过程的完整信息,保证可复现性
  5. 超参数调优系统性地搜索最优配置

五、特征工程进阶

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 系统中,手动创建所有特征是不现实的。自动化特征工程工具可以系统性地探索特征空间:

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 都提供了自动化的血缘追踪能力。

延伸阅读