Skip to content

自动化工作流

机器学习 - 使用 Pipelines 自动化工作流

Section titled “机器学习 - 使用 Pipelines 自动化工作流”

机器学习项目通常涉及一系列标准步骤:数据加载、预处理、特征工程、模型训练和评估。手动重复执行这些步骤可能很繁琐、容易出错,并且难以管理实验。Scikit-learn 的 Pipeline 对象提供了一种将多个处理步骤和最终估计器(如分类器或回归器)串联成一个对象的方式。

可以将 Pipeline 视为数据的流水线:原始数据从一端进入,经过各种转换阶段,训练好的模型或预测结果从另一端输出。这可以简化工作流程,改进代码组织,最重要的是有助于防止数据泄露(Data Leakage)。

一个典型的机器学习 Pipeline 可能包含以下逻辑块:

  • 数据摄取 (Data Ingestion): 加载原始数据。
  • 数据准备/预处理 (Data Preparation/Preprocessing): 数据清洗、处理缺失值、编码类别特征、缩放数值特征。
  • 特征工程/选择 (Feature Engineering/Selection): 创建新特征、选择相关特征(例如,使用 PCA 或统计检验)。
  • 模型训练 (Model Training): 将选定的机器学习算法拟合到准备好的数据上。
  • 模型评估 (Model Evaluation): 使用适当的指标评估模型在未见数据上的性能。
  • (可选) 模型部署 (Model Deployment): 使训练好的模型可用于预测。
  • (可选) 监控与再训练 (Monitoring & Retraining): 跟踪生产环境中的性能并按需进行再训练。

Scikit-learn 的 Pipeline 特别有助于自动化一系列转换器 (transformations) 后跟一个最终的估计器 (estimator) 的流程。

  • 便捷性与封装性: 将多个步骤组合到一个接口中。你只需在 Pipeline 对象上调用一次 fit 和 predict(或 transform)。
  • 防止数据泄露 (Data Leakage): 这是一个主要优势。Pipeline 确保预处理步骤(如数据缩放或填充缺失值)仅在交叉验证的每一折中的训练数据上进行拟合。在分割数据集之前对整个数据集应用预处理,可能因为测试集的信息“泄露”到训练过程中,导致模型性能评估过于乐观。
  • 联合参数选择: 允许同时对 Pipeline 中所有步骤的参数进行网格搜索 (grid search) 或随机搜索 (randomized search)。
  • 可复现性与可读性: 使操作序列清晰明确,更容易复现和理解。

尽管 Pipeline 有助于管理复杂性,但构建稳健的机器学习系统仍面临挑战:

  • 数据质量: Pipeline 依赖于输入数据。低质量(不准确、不一致、有偏见)的数据会导致模型性能低下(即“垃圾进,垃圾出”,Garbage In, Garbage Out)。
  • 数据漂移 (Data Drift)/概念漂移 (Concept Drift): 底层数据分布可能随时间变化,导致模型性能下降。Pipeline 需要监控和再训练机制。
  • 可扩展性: 高效处理大型数据集需要可扩展的工具和基础设施。
  • 复杂性: 设计包含许多步骤的复杂 Pipeline 需要仔细规划和验证。
  • 可解释性: 理解 Pipeline 做出某些预测的原因可能具有挑战性。
  • 偏差与公平性: 在有偏差的数据上训练的模型可能延续甚至放大社会偏差,导致不公平的结果。
  • 部署与监控 (MLOps): 将模型从研究阶段推向健壮的生产系统并随时间监控其性能,带来了工程挑战。
  • 专业知识需求: 设计、构建和维护有效的机器学习系统需要专业技能。

一个常见的用例是将预处理步骤(如数据缩放)与模型结合起来。Pipeline 确保缩放器 (scaler) 在交叉验证期间仅在数据的训练部分上进行拟合,从而防止数据泄露。

这个 Python 示例演示了如何创建一个 Pipeline,它首先使用 StandardScaler 对数据进行标准化,然后应用线性判别分析(Linear Discriminant Analysis, LDA)分类器。我们使用 K 折交叉验证(K-Fold cross-validation)在 Pima Indians Diabetes 数据集上评估整个 Pipeline。

import pandas as pd
from sklearn.model_selection import KFold, cross_val_score
from sklearn.preprocessing import StandardScaler
from sklearn.pipeline import Pipeline, make_pipeline # make_pipeline is often simpler
from sklearn.discriminant_analysis import LinearDiscriminantAnalysis
# Load data (replace path)
path = 'pima-indians-diabetes.csv'
names = ['preg', 'plas', 'pres', 'skin', 'test', 'mass', 'pedi', 'age', 'class']
dataframe = pd.read_csv(path, names=names)
# Separate features (X) and target (y)
array = dataframe.values
X = array[:, :-1]
y = array[:, -1]
# --- Create the Pipeline ---
# Method 1: Using Pipeline constructor (requires naming steps)
estimators = []
estimators.append(('standardize', StandardScaler()))
estimators.append(('lda', LinearDiscriminantAnalysis()))
model_pipeline = Pipeline(estimators)
# Method 2: Using make_pipeline (simpler, automatically names steps)
# model_pipeline = make_pipeline(StandardScaler(), LinearDiscriminantAnalysis())
# --- Evaluate the Pipeline ---
kf = KFold(n_splits=10, shuffle=True, random_state=42) # Use shuffle for better robustness
# cross_val_score handles fitting and scoring internally for each fold
# It ensures StandardScaler is fit only on the training fold each time
results = cross_val_score(model_pipeline, X, y, cv=kf, scoring='accuracy')
print(f"Pipeline Cross-Validation Accuracy: {results.mean():.4f} (+/- {results.std():.4f})")

输出结果显示了通过 10 折交叉验证获得的平均准确率 (accuracy)(和标准差)。这个分数反映了整个工作流(标准化后进行 LDA 分类)在稳健评估下的性能,防止了在缩放步骤中的数据泄露。

Pipeline 还可以包含特征提取或选择步骤。FeatureUnion 允许并行应用多个转换器 (transformer) 对象的输出结果进行组合。一个更现代化且通常更受欢迎的替代方案是 ColumnTransformer,它允许对输入数据的不同列应用不同的转换。

这个示例展示了一个更复杂的 Pipeline:

  1. 应用主成分分析(Principal Component Analysis, PCA)来降低维度。
  2. 应用单变量特征选择 (SelectKBest) 基于统计检验选择最佳特征。
  3. 使用 FeatureUnion 来连接 PCA 和 SelectKBest 的结果(注意:ColumnTransformer 通常更适合对不同的原始列应用不同的转换,而 FeatureUnion 适用于对相同的输入应用不同的转换并将结果连接起来)。
  4. 将组合后的特征输入到逻辑回归 (Logistic Regression) 分类器中。

整个组合过程使用交叉验证进行评估。

import pandas as pd
from sklearn.model_selection import KFold, cross_val_score
from sklearn.pipeline import Pipeline, FeatureUnion
from sklearn.linear_model import LogisticRegression
from sklearn.decomposition import PCA
from sklearn.feature_selection import SelectKBest, f_classif # Use f_classif for classification tasks
# Load data (replace path)
path = 'pima-indians-diabetes.csv'
names = ['preg', 'plas', 'pres', 'skin', 'test', 'mass', 'pedi', 'age', 'class']
dataframe = pd.read_csv(path, names=names)
# Separate features (X) and target (y)
array = dataframe.values
X = array[:, :-1]
y = array[:, -1]
# --- Create Feature Union ---
# Apply different feature extraction methods in parallel
features = []
features.append(('pca', PCA(n_components=3))) # 提取 3 个主成分
features.append(('select_best', SelectKBest(score_func=f_classif, k=6))) # 基于 ANOVA F-value 选择前 6 个特征
feature_union = FeatureUnion(features)
# --- Create the Full Pipeline ---
estimationers = []
estimationers.append(('feature_processing', feature_union)) # 步骤 1: 并行特征提取/选择
estimationers.append(('logistic', LogisticRegression(solver='liblinear', max_iter=200))) # 步骤 2: 分类器
model_pipeline = Pipeline(estimationers)
# --- Evaluate the Pipeline ---
kf = KFold(n_splits=10, shuffle=True, random_state=42)
results = cross_val_score(model_pipeline, X, y, cv=kf, scoring='accuracy')
print(f"Pipeline (PCA+SelectKBest -> LogisticRegression) Accuracy: {results.mean():.4f} (+/- {results.std():.4f})")
# --- 关于 ColumnTransformer 的说明(现代替代方案) ---
# 如果想对数值列应用 StandardScaler,对类别列应用 OneHotEncoder,
# ColumnTransformer 是合适的工具:
# from sklearn.compose import ColumnTransformer
# from sklearn.preprocessing import StandardScaler, OneHotEncoder
# preprocessor = ColumnTransformer(
# transformers=[
# ('num', StandardScaler(), numerical_features_indices),
# ('cat', OneHotEncoder(), categorical_features_indices)])
# pipe = Pipeline(steps=[('preprocessor', preprocessor), ('classifier', LogisticRegression())])

输出结果提供了整个复杂工作流的交叉验证准确率 (accuracy),包括并行特征处理和随后的分类。这展示了 Pipeline 如何封装复杂的、多步骤的流程,以实现稳健的评估和部署。

更多资源: