【问题标题】:Efficient way of row/column sum of a IndexedRowmatrix in Apache SparkApache Spark中IndexedRowmatrix的行/列总和的有效方法
【发布时间】:2016-01-23 04:21:16
【问题描述】:

我在 Scala 中有一个 CoordinateMatrix 格式的矩阵。矩阵是稀疏的,整体看起来像(根据 coo_matrix.entries.collect),

Array[org.apache.spark.mllib.linalg.distributed.MatrixEntry] = Array(
  MatrixEntry(0,0,-1.0), MatrixEntry(0,1,-1.0), MatrixEntry(1,0,-1.0),
  MatrixEntry(1,1,-1.0), MatrixEntry(1,2,-1.0), MatrixEntry(2,1,-1.0), 
  MatrixEntry(2,2,-1.0), MatrixEntry(0,3,-1.0), MatrixEntry(0,4,-1.0), 
  MatrixEntry(0,5,-1.0), MatrixEntry(3,0,-1.0), MatrixEntry(4,0,-1.0), 
  MatrixEntry(3,3,-1.0), MatrixEntry(3,4,-1.0), MatrixEntry(4,3,-1.0),
  MatrixEntry(4,4,-1.0))

这只是一个小样本。矩阵的大小为 N x N(其中 N = 100 万),尽管其中大部分是稀疏的。在 Spark Scala 中获取该矩阵的行和的有效方法之一是什么?目标是创建一个由行总和组成的新 RDD,即大小为 N,其中第一个元素是 row1 的行总和,依此类推..

我总是可以将此坐标矩阵转换为 IndexedRowMatrix 并运行一个 for 循环来一次计算一次迭代的行和,但这不是最有效的方法。

非常感谢任何想法。

【问题讨论】:

    标签: scala matrix apache-spark apache-spark-mllib rowsum


    【解决方案1】:

    由于洗牌,这将非常昂贵(这是您在此处无法真正避免的部分),但您可以将条目转换为 PairRDD 并按键减少:

    import org.apache.spark.mllib.linalg.distributed.{MatrixEntry, CoordinateMatrix}
    import org.apache.spark.rdd.RDD
    
    val mat: CoordinateMatrix = ???
    val rowSums: RDD[Long, Double)] = mat.entries
      .map{case MatrixEntry(row, _, value) => (row, value)}
      .reduceByKey(_ + _)
    

    不同于基于indexedRowMatrix的解决方案:

    import org.apache.spark.mllib.linalg.distributed.IndexedRow
    
    mat.toIndexedRowMatrix.rows.map{
      case IndexedRow(i, values) => (i, values.toArray.sum)
    }
    

    它不需要groupBy 转换或中间SparseVectors

    【讨论】:

      猜你喜欢
      • 2015-08-01
      • 1970-01-01
      • 2015-04-05
      • 1970-01-01
      • 2017-09-04
      • 2023-03-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多