【问题标题】:How should I convert an RDD of org.apache.spark.ml.linalg.Vector to Dataset?我应该如何将 org.apache.spark.ml.linalg.Vector 的 RDD 转换为数据集?
【发布时间】:2018-11-14 08:30:19
【问题描述】:

我很难理解 RDD、DataSet 和 DataFrame 之间的转换是如何工作的。 我对 Spark 还很陌生,每次我需要从一个数据模型传递到另一个(尤其是从 RDD 到数据集和数据帧)时,我都会卡住。 谁能解释一下正确的做法?

例如,现在我有一个RDD[org.apache.spark.ml.linalg.Vector],我需要将它传递给我的机器学习算法,例如 KMeans (Spark DataSet MLlib)。因此,我需要将其转换为具有名为“特征”的单列的数据集,该列应包含向量类型的行。我该怎么做?

【问题讨论】:

    标签: apache-spark apache-spark-sql rdd apache-spark-mllib apache-spark-dataset


    【解决方案1】:

    要将 RDD 转换为 数据帧,最简单的方法是在 Scala 中使用 toDF()。要使用此功能,必须导入使用 SparkSession 对象完成的隐式。可以这样做:

    val spark = SparkSession.builder().getOrCreate()
    import spark.implicits._
    
    val df = rdd.toDF("features")
    

    toDF() 接受元组的 RDD。当 RDD 由普通的 Scala 对象组成时,它们将被隐式转换,即不需要做任何事情,当 RDD 有多个列时也不需要做任何事情,RDD 已经包含一个元组。但是,在这种特殊情况中,您需要先将RDD[org.apache.spark.ml.linalg.Vector] 转换为RDD[(org.apache.spark.ml.linalg.Vector)]。因此,需要对元组进行如下转换:

    val df = rdd.map(Tuple1(_)).toDF("features")
    

    以上将 RDD 转换为具有称为特征的单列的数据框。


    要转换为数据集,最简单的方法是使用案例类。确保案例类是在 Main 对象之外定义的。先将RDD转成dataframe,然后进行如下操作:

    case class A(features: org.apache.spark.ml.linalg.Vector)
    
    val ds = df.as[A]
    

    要显示所有可能的转换,可以使用 .rdd 来从数据框或数据集访问底层 RDD

    val rdd = df.rdd
    

    与在 RDD 和数据帧/数据集之间来回转换不同,使用数据帧 API 进行所有计算通常更容易。如果没有合适的函数来做你想做的事,通常可以定义一个 UDF,用户定义的函数。例如,请参见此处:https://jaceklaskowski.gitbooks.io/mastering-spark-sql/spark-sql-udfs.html

    【讨论】:

    • 我试过你的解决方案,它给了我这个错误:错误:(94, 32) value toDF is not a member of org.apache.spark.rdd.RDD[org.apache.spark .ml.linalg.Vector] val df = rdd.toDF("features")
    • @Giuseppe 你好像忘了import sqlContext.implicits._
    • 我尝试了一个和两个import sqlContext.implicits._ import spark.implicits._ 但没有任何改变。我在主目录中创建了上下文和会话之后才将它们放入。会不会是spark版本的问题?我正在使用 Spark 2.2.1
    • 如果它可以帮助发现问题,如果我转换为元组rdd.map(x => ("",x)).toDF() toDF() 工作。此时我只需要投影 _2 列并重命名它。
    • @Giuseppe:这有帮助,问题是由于 RDD 需要一个元组。我在答案中添加了更多信息,应该有助于解决问题:)
    【解决方案2】:

    您只需要一个Encoder。进口

    import org.apache.spark.sql.Encoder
    import org.apache.spark.sql.catalyst.encoders.ExpressionEncoder
    import org.apache.spark.ml.linalg
    

    RDD:

    val rdd = sc.parallelize(Seq(
      linalg.Vectors.dense(1.0, 2.0), linalg.Vectors.sparse(2, Array(), Array())
    ))
    

    转化率:

    val ds = spark.createDataset(rdd)(ExpressionEncoder(): Encoder[linalg.Vector])
     .toDF("features")
    
    ds.show
    // +---------+
    // | features|
    // +---------+
    // |[1.0,2.0]|
    // |(2,[],[])|
    // +---------+
    
    
    ds.printSchema
    // root
    //  |-- features: vector (nullable = true)
    

    【讨论】:

      猜你喜欢
      • 2017-07-08
      • 2018-06-14
      • 2016-12-12
      • 1970-01-01
      • 1970-01-01
      • 2020-01-24
      • 1970-01-01
      • 1970-01-01
      • 2017-04-08
      相关资源
      最近更新 更多