【问题标题】:Save Spark org.apache.spark.mllib.linalg.Matrix to a file将 Spark org.apache.spark.mllib.linalg.Matrix 保存到文件中
【发布时间】:2017-05-19 14:27:21
【问题描述】:

Spark MLLib 中的关联结果是 org.apache.spark.mllib.linalg.Matrix 类型。 (见http://spark.apache.org/docs/1.2.1/mllib-statistics.html#correlations

val data: RDD[Vector] = ... 

val correlMatrix: Matrix = Statistics.corr(data, "pearson")

我想将结果保存到文件中。我该怎么做?

【问题讨论】:

    标签: apache-spark apache-spark-mllib


    【解决方案1】:

    这里有一个简单有效的方法,将Matrix保存到hdfs并指定分隔符。

    (使用转置,因为 .toArray 是列主要格式。)

    val localMatrix: List[Array[Double]] = correlMatrix
        .transpose  // Transpose since .toArray is column major
        .toArray
        .grouped(correlMatrix.numCols)
        .toList
    
    val lines: List[String] = localMatrix
        .map(line => line.mkString(" "))
    
    sc.parallelize(lines)
        .repartition(1)
        .saveAsTextFile("hdfs:///home/user/spark/correlMatrix.txt")
    

    【讨论】:

      【解决方案2】:

      由于Matrix 是可序列化的,您可以使用普通的Scala 编写它。

      你可以找到一个例子here

      【讨论】:

      • 感谢您的回答,卡洛斯。我想将矩阵保存在 HDFS 中。此外,如果可能的话,采用人类“可读”的格式。就像是。 RDD API 提供的 saveAsTextFile
      • 你可以试试data.saveAsTextFile("hdfs://...")。我在Spark examples web看到过。
      【解决方案3】:

      Dylan Hogg 的回答很棒,为了稍微增强它,添加一个列索引。 (在我的用例中,一旦我创建了一个文件并下载了它,由于并行过程等的性质,它没有被排序。)

      参考:https://www.safaribooksonline.com/library/view/scala-cookbook/9781449340292/ch10s12.html

      替换为这一行,它会在该行上放置一个序列号(从 0 开始),以便您查看时更容易排序

      val lines: List[String] = localMatrix 
        .map(line => line.mkString(" ")) 
        .zipWithIndex.map { case(line, count) => s"$count $line" } 
      

      【讨论】:

        【解决方案4】:

        感谢您的建议。我提出了这个解决方案。感谢 Ignacio 的建议

        val vtsd = sd.map(x => Vectors.dense(x.toArray))
        val corrMat = Statistics.corr(vtsd)
        val arrayCor = corrMat.toArray.toList
        val colLen = columnHeader.size
        val toArr2 = sc.parallelize(arrayCor).zipWithIndex().map(
              x => {
            if ((x._2 + 1) % colLen == 0) {
              (x._2, arrayCor.slice(x._2.toInt + 1 - colLen, x._2.toInt + 1).mkString(";"))
            } else {
              (x._2, "")
            }
          }).filter(_._2.nonEmpty).sortBy(x => x._1, true, 1).map(x => x._2)
        
        
        toArr2.coalesce(1, true).saveAsTextFile("/home/user/spark/cor_" + System.currentTimeMillis())
        

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 2018-12-05
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2017-09-12
          • 2020-09-15
          • 1970-01-01
          相关资源
          最近更新 更多