【问题标题】:Convert Rdd[Vector] to Rdd[Double]将 Rdd[Vector] 转换为 Rdd[Double]
【发布时间】:2015-10-05 09:28:00
【问题描述】:

如何将 csv 转换为 Rdd[Double]?我有错误:无法在此行应用于 (org.apache.spark.rdd.RDD[Unit]):

val kd = new KernelDensity().setSample(rows) 

我的完整代码在这里:

   import org.apache.spark.mllib.linalg.Vectors
    import org.apache.spark.mllib.linalg.distributed.RowMatrix
    import org.apache.spark.mllib.stat.KernelDensity
    import org.apache.spark.rdd.RDD
    import org.apache.spark.{SparkContext, SparkConf}

class KdeAnalysis {
  val conf = new SparkConf().setAppName("sample").setMaster("local")
  val sc = new SparkContext(conf)

  val DATAFILE: String = "C:\\Users\\ajohn\\Desktop\\spark_R\\data\\mass_cytometry\\mass.csv"
  val rows = sc.textFile(DATAFILE).map {
    line => val values = line.split(',').map(_.toDouble)
      Vectors.dense(values)
  }.cache()



  // Construct the density estimator with the sample data and a standard deviation for the Gaussian
  // kernels
  val rdd : RDD[Double] = sc.parallelize(rows)
  val kd = new KernelDensity().setSample(rdd)
    .setBandwidth(3.0)

  // Find density estimates for the given values
  val densities = kd.estimate(Array(-1.0, 2.0, 5.0))
}

【问题讨论】:

  • 我没有看到你的任何地方可以得到RDD[Unit]

标签: scala apache-spark rdd apache-spark-mllib


【解决方案1】:

由于rowsRDD[org.apache.spark.mllib.linalg.Vector],以下行无法工作:

val rdd : RDD[Double] = sc.parallelize(rows)

parallelize 期望 Seq[T]RDD 不是 Seq

即使这部分按照您的预期工作,您的输入也完全是错误的。 KernelDensity.setSample 的正确参数是 RDD[Double]JavaRDD[java.lang.Double]。目前看来它不支持多变量数据。

关于磁贴的问题,您可以flatMap

rows.flatMap(_.toArray)

创建rows时甚至更好

val rows = sc.textFile(DATAFILE).flatMap(_.split(',').map(_.toDouble)).cache()

但我怀疑它真的是你需要的。

【讨论】:

    【解决方案2】:

    已准备好此代码,请评估它是否可以帮助您->

    val doubleRDD = rows.map(_.toArray).flatMap(x => x)
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2021-11-30
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-06-29
      • 2017-06-02
      • 2016-01-07
      • 2021-02-26
      相关资源
      最近更新 更多