流水线

使用 PySpark 进行机器学习

Andrew Collier

Data Scientist, Fathom Data

数据泄漏?

fit() 方法

仅用于训练数据。

transform() 方法

用于测试和训练数据。

使用 PySpark 进行机器学习

有泄漏的模型

将测试数据用于训练的模型

使用 PySpark 进行机器学习

无泄漏的模型

仅用训练数据进行训练的模型

使用 PySpark 进行机器学习

流水线

流水线由一系列操作组成。

包含多个阶段的流水线

可以逐个执行每个操作,也可以直接运行整条流水线!

使用 PySpark 进行机器学习

汽车模型:步骤

indexer = StringIndexer(inputCol='type', outputCol='type_idx')

onehot = OneHotEncoder(inputCols=['type_idx'], outputCols=['type_dummy'])
assemble = VectorAssembler( inputCols=['mass', 'cyl', 'type_dummy'], outputCol='features' )
regression = LinearRegression(labelCol='consumption')
使用 PySpark 进行机器学习

汽车模型:应用各步骤

训练数据

indexer = indexer.fit(cars_train)
cars_train = indexer.transform(cars_train)
onehot = onehot.fit(cars_train)
cars_train = onehot.transform(cars_train)
cars_train = assemble.transform(cars_train)
# 在训练数据上拟合模型
regression = regression.fit(cars_train)

测试数据

cars_test  = indexer.transform(cars_test)
cars_test  = onehot.transform(cars_test)
cars_test  = assemble.transform(cars_test)
# 在测试数据上生成预测
predictions = regression.transform(cars_test)
使用 PySpark 进行机器学习

汽车模型:流水线

将步骤合并为一个流水线。

from pyspark.ml import Pipeline

pipeline = Pipeline(stages=[indexer, onehot, assemble, regression])

训练数据

pipeline = pipeline.fit(cars_train)

测试数据

predictions = pipeline.transform(cars_test)
使用 PySpark 进行机器学习

汽车模型:各阶段

使用 .stages 属性访问各阶段。

# 线性回归对象(第4阶段 -> 索引3)
pipeline.stages[3]

print(pipeline.stages[3].intercept)
4.19433571782916
print(pipeline.stages[3].coefficients)
DenseVector([0.0028, 0.2705, -1.1813, -1.3696, -1.1751, -1.1553, -1.8894])
使用 PySpark 进行机器学习

流水线让流程更顺畅!

使用 PySpark 进行机器学习

Preparing Video For Download...