【问题标题】:Addition of two RDD[mllib.linalg.Vector]'s添加两个 RDD[mllib.linalg.Vector]
【发布时间】:2014-09-27 21:52:39
【问题描述】:

我需要添加两个存储在两个文件中的矩阵。

latest1.txtlatest2.txt的内容有下一个str:

1 2 3 4 5 6 7 8 9

我正在阅读这些文件如下:

scala> val rows = sc.textFile(“latest1.txt”).map { line => val values = line.split(‘ ‘).map(_.toDouble)
    Vectors.sparse(values.length,values.zipWithIndex.map(e => (e._2, e._1)).filter(_._2 != 0.0))
}

scala> val r1 = rows
r1: org.apache.spark.rdd.RDD[org.apache.spark.mllib.linalg.Vector] = MappedRDD[2] at map at :14

scala> val rows = sc.textFile(“latest2.txt”).map { line => val values = line.split(‘ ‘).map(_.toDouble)
    Vectors.sparse(values.length,values.zipWithIndex.map(e => (e._2, e._1)).filter(_._2 != 0.0))
}

scala> val r2 = rows
r2: org.apache.spark.rdd.RDD[org.apache.spark.mllib.linalg.Vector] = MappedRDD[2] at map at :14

我想添加 r1、r2。那么,有没有办法在 Apache-Spark 中添加这两个 RDD[mllib.linalg.Vector]s。

【问题讨论】:

  • 将两个 RDD 压缩在一起,然后映射到生成的 RDD
  • 是的,我确实喜欢 val rdd3=rdd1.zip(rdd2) scala> val rdd4 = rdd3.map{ e => e._1 + e._2} 我收到错误:22:错误:类型不匹配;找到: org.apache.spark.mllib.linalg.Vector 需要:String val r4=r3.map{e=>e._1 + e._2} 因为在 mllib 向量上没有 + 或 add 操作,所以定义了加法操作util.Vectors
  • 看起来 + 不是添加两个向量的运算符,所以你得到了尝试转换为字符串的默认隐含。
  • 是的,但我找不到任何执行加法的函数或运算符。

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


【解决方案1】:

这实际上是一个很好的问题。我经常使用 mllib,但没有意识到这些基本的线性代数运算不容易访问。

关键在于,底层的 breeze 向量具有您所期望的所有线性代数操作 - 当然包括您特别提到的基本元素加法。

然而,微风的实现是通过以下方式对外界隐藏的:

[private mllib]

那么,从外部世界/公共 API 的角度来看,我们如何访问这些原语?

其中一些已经暴露:例如平方和:

/**
 * Returns the squared distance between two Vectors.
 * @param v1 first Vector.
 * @param v2 second Vector.
 * @return squared distance between two Vectors.
 */
def sqdist(v1: Vector, v2: Vector): Double = { 
  ...
}

然而,这些可用方法的选择是有限的——事实上包括基本操作,包括元素加法、减法、乘法等。

所以这是我能看到的最好的:

  • 将向量转换为微风:
  • 在微风中执行矢量运算
  • 从微风转换回 mllib 向量

这里是一些示例代码:

val v1 = Vectors.dense(1.0, 2.0, 3.0)
val v2 = Vectors.dense(4.0, 5.0, 6.0)
val bv1 = new DenseVector(v1.toArray)
val bv2 = new DenseVector(v2.toArray)

val vectout = Vectors.dense((bv1 + bv2).toArray)
vectout: org.apache.spark.mllib.linalg.Vector = [5.0,7.0,9.0]

【讨论】:

  • 是的。 MLlib 不是一个完整的线性代数库,如果需要这样的操作,应该使用Breeze
  • 但是如果向量是稀疏的怎么办。我目前正在处理稀疏向量。但是如果使用你的方式转换向量,会消耗更多的内存并降低计算速度。奇怪的是,pyspark 可以轻松完成此操作。所以我正在考虑改用python。
  • 这实际上是我正在尝试的.. but what am I doing wrong here?
  • @displayname 我回答了这个问题。
  • @javadba 您认为处理稀疏向量时性能会受到多大影响?我正在处理长度为 2**20 的 Spark 向量,但我似乎无法在 Scala 中找到一种有效的方法来处理这个问题。
【解决方案2】:

以下代码公开了 Spark 的 asBreeze 和 fromBreeze 方法。与使用vector.toArray 相比,此解决方案支持SparseVector。请注意,Spark 将来可能会更改其 API,并且已经将 toBreeze 重命名为 asBreeze

package org.apache.spark.mllib.linalg
import breeze.linalg.{Vector => BV}
import org.apache.spark.sql.functions.udf

/** expose vector.toBreeze and Vectors.fromBreeze
  */
object VectorUtils {

  def fromBreeze(breezeVector: BV[Double]): Vector = {
    Vectors.fromBreeze( breezeVector )
  }

  def asBreeze(vector: Vector): BV[Double] = {
    // this is vector.asBreeze in Spark 2.0
    vector.toBreeze
  }

  val addVectors = udf {
    (v1: Vector, v2: Vector) => fromBreeze( asBreeze(v1) + asBreeze(v2) )
  }

}

有了这个你可以df.withColumn("xy", addVectors($"x", $"y"))

【讨论】:

  • 第一行不应该是import org.apache.spark.mllib.linalg._而不是包定义吗?如果我按原样使用,我会收到一条错误消息“非法开始定义”。
  • @scottH 否,因为函数需要成为包的一部分才能访问私有函数。代码在 Spark 1.6.1 中运行良好,但 Spark 2+ 改变了一切。您是否尝试将代码编译成 JAR 而不是复制粘贴到 spark-shell?
猜你喜欢
  • 1970-01-01
  • 2016-09-07
  • 2016-01-24
  • 1970-01-01
  • 2023-03-13
  • 1970-01-01
  • 2017-07-30
  • 2015-06-15
相关资源
最近更新 更多