使用 PySpark 进行机器学习
Andrew Collier
Data Scientist, Fathom Data
仅用于训练数据。
用于测试和训练数据。


流水线由一系列操作组成。
可以逐个执行每个操作,也可以直接运行整条流水线!
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')
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)
将步骤合并为一个流水线。
from pyspark.ml import Pipeline
pipeline = Pipeline(stages=[indexer, onehot, assemble, regression])
训练数据
pipeline = pipeline.fit(cars_train)
测试数据
predictions = pipeline.transform(cars_test)
使用 .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 进行机器学习