PySpark PCA:如何将数据帧行从多列转换为单列DenseVector?

Ale*_*ord 4 pca apache-spark pyspark apache-spark-ml apache-spark-mllib

我想使用PySpark(Spark 1.6.2)对Hive表中存在的数值数据执行主成分分析(PCA).我能够将Hive表导入Spark数据帧:

>>> from pyspark.sql import HiveContext
>>> hiveContext = HiveContext(sc)
>>> dataframe = hiveContext.sql("SELECT * FROM my_table")
>>> type(dataframe)
<class 'pyspark.sql.dataframe.DataFrame'>
>>> dataframe.columns
['par001', 'par002', 'par003', etc...]
>>> dataframe.collect()
[Row(par001=1.1, par002=5.5, par003=8.2, etc...), Row(par001=0.0, par002=5.7, par003=4.2, etc...), etc...]
Run Code Online (Sandbox Code Playgroud)

有一个很棒的StackOverflow帖子,展示了如何在PySpark中执行PCA:https://stackoverflow.com/a/33481471/2626491

在帖子的"测试"部分,@ assellnaut创建了一个只有一列的数据框(称为"要素"):

>>> from pyspark.ml.feature import *
>>> from pyspark.mllib.linalg import Vectors
>>> data = [(Vectors.dense([0.0, 1.0, 0.0, 7.0, 0.0]),),
...          (Vectors.dense([2.0, 0.0, 3.0, 4.0, 5.0]),),
...          (Vectors.dense([4.0, 0.0, 0.0, 6.0, 7.0]),)]
>>> df = sqlContext.createDataFrame(data,["features"])
>>> type(df)
<class 'pyspark.sql.dataframe.DataFrame'>
>>> df.columns
['features']
>>> df.collect()
[Row(features=DenseVector([0.0, 1.0, 0.0, 7.0, 0.0])), Row(features=DenseVector([2.0, 0.0, 3.0, 4.0, 5.0])), Row(features=DenseVector([4.0, 0.0, 0.0, 6.0, 7.0]))]
Run Code Online (Sandbox Code Playgroud)

@ desertnaut的示例数据框中的每一行都包含一个DenseVector对象,然后由该pca函数使用.

问)如何将数据框从Hive转换为单列数据框("features"),其中每行包含一个DenseVector代表原始行的所有值?

use*_*411 8

你应该用一个VectorAssembler.如果数据与此类似:

from pyspark.sql import Row

data = sc.parallelize([
    Row(par001=1.1, par002=5.5, par003=8.2),
    Row(par001=0.0, par002=5.7, par003=4.2)
]).toDF()
Run Code Online (Sandbox Code Playgroud)

你应该导入所需的类:

from pyspark.ml.feature import VectorAssembler
Run Code Online (Sandbox Code Playgroud)

创建一个实例:

assembler = VectorAssembler(inputCols=data.columns, outputCol="features")
Run Code Online (Sandbox Code Playgroud)

转换并选择:

assembler.transform(data).select("features")
Run Code Online (Sandbox Code Playgroud)

您还可以使用用户定义的函数.在Spark 1.6导入VectorsVectorUDT来自mllib:

from pyspark.mllib.linalg import Vectors, VectorUDT
Run Code Online (Sandbox Code Playgroud)

udf来自sql.functions:

from pyspark.sql.functions import udf, array
Run Code Online (Sandbox Code Playgroud)

并选择:

data.select(
  udf(Vectors.dense, VectorUDT())(*data.columns)
).toDF("features")
Run Code Online (Sandbox Code Playgroud)

这不是那么冗长,而是慢得多.