【问题标题】:Create labledpoints from Spark Dataframe & how to pass list of names to VectorAssembler从 Spark Dataframe 创建标签点以及如何将名称列表传递给 VectorAssembler
【发布时间】:2019-01-13 00:13:25
【问题描述】:

我还有其他问题要从这里https://stackoverflow.com/a/32557330/5235052 我正在尝试从数据框中构建 labeldPoints,其中我有列中的特征和标签。这些特征都是布尔值,带 1/0。

这是数据框中的示例行:

|             0|       0|        0|            0|       0|            0|     1|        0|     0|           0|       0|       0|       0|           0|        0|         0|      0|            0|       0|           0|          0|         0|         0|              0|        0|        0|        0|         0|          0|    1|    0|    1|    0|    0|       0|           0|    0|     0|     0|     0|         0|         1|
#Using the code from above answer, 
#create a list of feature names from the column names of the dataframe
df_columns = []
for  c in df.columns:
    if c == 'is_item_return': continue
    df_columns.append(c)

#using VectorAssembler for transformation, am using only first 4 columns names
assembler = VectorAssembler()
assembler.setInputCols(df_columns[0:5])
assembler.setOutputCol('features')

transformed = assembler.transform(df)

   #mapping also from above link
   from pyspark.mllib.regression import LabeledPoint
   from pyspark.sql.functions import col

new_df = transformed.select(col('is_item_return'), col("features")).map(lambda row: LabeledPoint(row.is_item_return, row.features))

当我检查 RDD 的内容时,我得到了正确的标签,但特征向量是错误的。

(0.0,(5,[],[]))

有人可以帮助我理解,如何将现有数据框的列名作为特征名传递给 VectorAssembler?

【问题讨论】:

    标签: python apache-spark apache-spark-sql apache-spark-ml


    【解决方案1】:

    这里没有错。您得到的是 SparseVector 的字符串表示形式,它完全反映了您的输入:

    • 您取前五列 (assembler.setInputCols(df_columns[0:5])),输出向量的长度为 5
    • 由于示例输入的第一列不包含非零条目,indicesvalues 数组为空

    为了说明这一点,让我们使用提供有用的toSparse / toDense 方法的 Scala:

    import org.apache.spark.mllib.linalg.Vectors
    
    val v = Vectors.dense(Array(0.0, 0.0, 0.0, 0.0, 0.0))
    v.toSparse.toString
    // String = (5,[],[])
    
    v.toSparse.toDense.toString
    // String = [0.0,0.0,0.0,0.0,0.0]
    

    PySpark 也是如此:

    from pyspark.ml.feature import VectorAssembler
    
    df = sc.parallelize([
        tuple([0.0] * 5),
        tuple([1.0] * 5), 
        (1.0, 0.0, 1.0, 0.0, 1.0),
        (0.0, 1.0, 0.0, 1.0, 0.0)
    ]).toDF()
    
    features = (VectorAssembler(inputCols=df.columns, outputCol="features")
        .transform(df)
        .select("features"))
    
    features.show(4, False)
    
    ## +---------------------+
    ## |features             |
    ## +---------------------+
    ## |(5,[],[])            |
    ## |[1.0,1.0,1.0,1.0,1.0]|
    ## |[1.0,0.0,1.0,0.0,1.0]|
    ## |(5,[1,3],[1.0,1.0])  |
    ## +---------------------+
    

    它还表明汇编器根据非零条目的数量选择不同的表示。

    features.flatMap(lambda x: x).map(type).collect()
    
    ## [pyspark.mllib.linalg.SparseVector,
    ##  pyspark.mllib.linalg.DenseVector,
    ##  pyspark.mllib.linalg.DenseVector,
    ##  pyspark.mllib.linalg.SparseVector]
    

    【讨论】:

    • 这很清楚,谢谢!。来自 pandas-scikit 的新火花,有一个后续问题,我如何在这个 rdd 上运行线性回归? * 我是创建一个数据框然后使用 ML 还是将 rdd 保持这种格式并使用 MLLib? * 我所指的例子对理解输入数据结构没有多大帮助。link * 在 pandas 中,我会将数据拆分为独立的 x 和 y,创建训练、测试。拟合模型,测试分数得到系数。 * 如何在 spark ML/MLib 中做到这一点?
    • 更新,我可以使用from pyspark.ml.regression import LinearRegression 运行 lr 模型。我可以运行一个小型测试模型,但是当我将其扩展到 56 个功能、110 万行时,我的代码会崩溃。在distinctValues = df.map(lambda x : x[c]).distinct().collect() 这一行出现错误py4j.protocol.Py4JJavaError: An error occurred while calling z:org.apache.spark.api.python.PythonRDD.collectAndServe. : org.apache.spark.SparkException: Job aborted due to stage failure: Task serialization failed: java.lang.StackOverflowError 我在 Google 云上使用 Dataproc。
    • Mllib 回归已被弃用。关于内存问题,收集除小聚合之外的任何东西是没有意义的,一般来说是危险的。但这些都与这个问题无关,是吗?
    猜你喜欢
    • 2016-03-31
    • 2019-06-27
    • 1970-01-01
    • 2019-02-12
    • 1970-01-01
    • 2012-11-19
    • 1970-01-01
    • 1970-01-01
    • 2020-01-31
    相关资源
    最近更新 更多