使用 PySpark 进行机器学习
Andrew Collier
Data Scientist, Fathom Data
+-----+-------+-------+------+----+----+------+------+----+-----------+
|maker| model| origin| type| cyl|size|weight|length| rpm|consumption|
+-----+-------+-------+------+----+----+------+------+----+-----------+
|Mazda| RX-7|non-USA|Sporty|null| 1.3| 2895| 169.0|6500| 9.41|
| Geo| Metro|non-USA| Small| 3| 1.0| 1695| 151.0|5700| 4.7|
| Ford|Festiva| USA| Small| 4| 1.3| 1845| 141.0|5000| 7.13|
+-----+-------+-------+------+----+----+------+------+----+-----------+
移除 maker 和 model 字段。
# 要么丢弃不需要的列… cars = cars.drop('maker', 'model')# …要么仅选择保留的列。 cars = cars.select('origin', 'type', 'cyl', 'size', 'weight', 'length', 'rpm', 'consumption')
+-------+------+----+----+------+------+----+-----------+
| origin| type| cyl|size|weight|length| rpm|consumption|
+-------+------+----+----+------+------+----+-----------+
|non-USA|Sporty|null| 1.3| 2895| 169.0|6500| 9.41|
|non-USA| Small| 3| 1.0| 1695| 151.0|5700| 4.7|
| USA| Small| 4| 1.3| 1845| 141.0|5000| 7.13|
+-------+------+----+----+------+------+----+-----------+
# 有多少缺失值?
cars.filter('cyl IS NULL').count()
1
删除 cylinders 列中缺失值的记录。
cars = cars.filter('cyl IS NOT NULL')
删除任意列存在缺失值的记录。
cars = cars.dropna()
from pyspark.sql.functions import round # 新增"mass"列 cars = cars.withColumn('mass', round(cars.weight / 2.205, 0))# 将 length 转为米 cars = cars.withColumn('length', round(cars.length * 0.0254, 3))
+-------+-----+---+----+------+------+----+-----------+-----+
| origin| type|cyl|size|weight|length| rpm|consumption| mass|
+-------+-----+---+----+------+------+----+-----------+-----+
|non-USA|Small| 3| 1.0| 1695| 3.835|5700| 4.7|769.0|
| USA|Small| 4| 1.3| 1845| 3.581|5000| 7.13|837.0|
|non-USA|Small| 3| 1.3| 1965| 4.089|6000| 5.47|891.0|
+-------+-----+---+----+------+------+----+-----------+-----+
from pyspark.ml.feature import StringIndexer indexer = StringIndexer(inputCol='type', outputCol='type_idx')# 为字符串分配索引值 indexer = indexer.fit(cars)# 生成索引列 cars = indexer.transform(cars)
使用 stringOrderType 更改顺序。
+-------+--------+
| type|type_idx|
+-------+--------+
|Midsize| 0.0| <- most frequent value
| Small| 1.0|
|Compact| 2.0|
| Sporty| 3.0|
| Large| 4.0|
| Van| 5.0| <- least frequent value
+-------+--------+
# 为产地建立索引:
#
# USA -> 0
# non-USA -> 1
#
cars = StringIndexer(
inputCol="origin",
outputCol="label"
).fit(cars).transform(cars)
+-------+-----+
| origin|label|
+-------+-----+
| USA| 0.0|
|non-USA| 1.0|
+-------+-----+
使用向量组装器转换数据。
from pyspark.ml.feature import VectorAssembler assembler = VectorAssembler(inputCols=['cyl', 'size'], outputCol='features')assembler.transform(cars)
+---+----+---------+
|cyl|size| features|
+---+----+---------+
| 3| 1.0|[3.0,1.0]|
| 4| 1.3|[4.0,1.3]|
| 3| 1.3|[3.0,1.3]|
+---+----+---------+
使用 PySpark 进行机器学习