Man*_*aka 8 bigdata apache-spark distance-matrix
我已经尝试过对样本进行配对,但是它需要大量的内存,因为100个样本会导致9900个样本的成本更高.什么是在火花中分布式环境中计算距离矩阵的更有效方法
这是我正在尝试的伪代码片段
val input = (sc.textFile("AirPassengers.csv",(numPartitions/2)))
val i = input.map(s => (Vectors.dense(s.split(',').map(_.toDouble))))
val indexed = i.zipWithIndex() //Including the index of each sample
val indexedData = indexed.map{case (k,v) => (v,k)}
val pairedSamples = indexedData.cartesian(indexedData)
val filteredSamples = pairedSamples.filter{ case (x,y) =>
(x._1.toInt > y._1.toInt) //to consider only the upper or lower trainagle
}
filteredSamples.cache
filteredSamples.count
Run Code Online (Sandbox Code Playgroud)
上面的代码创建了对,但即使我的数据集包含100个样本,通过配对filteredSamples(上面)会产生4950样本,这对于大数据来说可能非常昂贵
我最近回答了一个类似的问题。
基本上,它将到达计算n(n-1)/2对,这将是4950您示例中的计算。但是,这种方法的不同之处在于我使用连接而不是cartesian. 使用您的代码,解决方案如下所示:
val input = (sc.textFile("AirPassengers.csv",(numPartitions/2)))
val i = input.map(s => (Vectors.dense(s.split(',').map(_.toDouble))))
val indexed = i.zipWithIndex()
// including the index of each sample
val indexedData = indexed.map { case (k,v) => (v,k) }
// prepare indices
val count = i.count
val indices = sc.parallelize(for(i <- 0L until count; j <- 0L until count; if i > j) yield (i, j))
val joined1 = indices.join(indexedData).map { case (i, (j, v)) => (j, (i,v)) }
val joined2 = joined1.join(indexedData).map { case (j, ((i,v1),v2)) => ((i,j),(v1,v2)) }
// after that, you can then compute the distance using your distFunc
val distRDD = joined2.mapValues{ case (v1, v2) => distFunc(v1, v2) }
Run Code Online (Sandbox Code Playgroud)
试试这个方法并将其与您已经发布的方法进行比较。希望这可以稍微加速您的代码。
小智 1
据我通过检查各种来源和Spark mllib 聚类站点发现,Spark 目前不支持距离或 pdist 矩阵。
在我看来,100 个样本总是会输出至少 4950 个值;因此,使用转换(如 .map)手动创建分布式矩阵求解器将是最好的解决方案。
| 归档时间: |
|
| 查看次数: |
2845 次 |
| 最近记录: |