【问题标题】:Spark Jaccard similarity computation by min hashing slow compared to trivial approachSpark Jaccard 相似度计算通过最小散列比普通方法慢
【发布时间】:2016-07-21 08:30:56
【问题描述】:

鉴于 2 个巨大的值列表,我正在尝试使用 Scala 在 Spark 中计算它们之间的 jaccard similarity。

假设colHashed1 包含第一个值列表,colHashed2 包含第二个列表。

方法一(普通方法):

val jSimilarity = colHashed1.intersection(colHashed2).distinct.count/(colHashed1.union(colHashed2).distinct.count.toDouble)

方法2(使用minHashing):

我使用了here解释的方法。

import java.util.zip.CRC32

def getCRC32 (s : String) : Int =
{
    val crc=new CRC32
    crc.update(s.getBytes)
    return crc.getValue.toInt & 0xffffffff
}

val maxShingleID = Math.pow(2,32)-1
def pickRandomCoeffs(kIn : Int) : Array[Int] =
{
  var k = kIn
  val randList = Array.fill(k){0}

  while(k > 0)
  {
    // Get a random shingle ID.

    var randIndex = (Math.random()*maxShingleID).toInt

    // Ensure that each random number is unique.
    while(randList.contains(randIndex))
    {
      randIndex = (Math.random()*maxShingleID).toInt
    }

    // Add the random number to the list.
    k = k - 1
    randList(k) = randIndex
   } 

   return randList
}

val colHashed1 = list1Values.map(a => getCRC32(a))
val colHashed2 = list2Values.map(a => getCRC32(a))

val nextPrime = 4294967311L
val numHashes = 10

val coeffA = pickRandomCoeffs(numHashes)
val coeffB = pickRandomCoeffs(numHashes)

var signature1 = Array.fill(numHashes){0}
for (i <- 0 to numHashes-1)
{
    // Evaluate the hash function.
    val hashCodeRDD = colHashed1.map(ele => ((coeffA(i) * ele + coeffB(i)) % nextPrime))

    // Track the lowest hash code seen.
    signature1(i) = hashCodeRDD.min.toInt
}

var signature2 = Array.fill(numHashes){0}
for (i <- 0 to numHashes-1)
{
    // Evaluate the hash function.
    val hashCodeRDD = colHashed2.map(ele => ((coeffA(i) * ele + coeffB(i)) % nextPrime))

    // Track the lowest hash code seen.
    signature2(i) = hashCodeRDD.min.toInt
}


var count = 0
// Count the number of positions in the minhash signature which are equal.
for(k <- 0 to numHashes-1)
{
  if(signature1(k) == signature2(k))
    count = count + 1
}  
val jSimilarity = count/numHashes.toDouble

方法 1 在时间方面似乎总是优于方法 2。当我分析代码时,在方法 2 中对 RDD 的 min() 函数调用需要大量时间,并且该函数被调用多次,具体取决于使用了多少哈希函数。

与重复的 min() 函数调用相比,方法 1 中使用的交集和并集操作似乎工作得更快。

我不明白为什么 minHashing 在这里没有帮助。与琐碎的方法相比,我希望 minHashing 工作得更快。我在这里做错了什么吗?

样本数据可以查看here

【问题讨论】:

  • 您可以在数据集中为您的 col1 和 col2 添加示例数据吗?
  • @tuxdna 示例数据链接添加在问题的末尾

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


【解决方案1】:

JaccardSimilarity 与 MinHash 的结果不一致:

import java.util.zip.CRC32

object Jaccard {
  def getCRC32(s: String): Int = {
    val crc = new CRC32
    crc.update(s.getBytes)
    return crc.getValue.toInt & 0xffffffff
  }

  def pickRandomCoeffs(kIn: Int, maxShingleID: Double): Array[Int] = {
    var k = kIn
    val randList = Array.ofDim[Int](k)

    while (k > 0) {
      // Get a random shingle ID.
      var randIndex = (Math.random() * maxShingleID).toInt
      // Ensure that each random number is unique.
      while (randList.contains(randIndex)) {
        randIndex = (Math.random() * maxShingleID).toInt
      }
      // Add the random number to the list.
      k = k - 1
      randList(k) = randIndex
    }
    return randList
  }


  def approach2(list1Values: List[String], list2Values: List[String]) = {

    val maxShingleID = Math.pow(2, 32) - 1

    val colHashed1 = list1Values.map(a => getCRC32(a))
    val colHashed2 = list2Values.map(a => getCRC32(a))

    val nextPrime = 4294967311L
    val numHashes = 10

    val coeffA = pickRandomCoeffs(numHashes, maxShingleID)
    val coeffB = pickRandomCoeffs(numHashes, maxShingleID)

    val signature1 = for (i <- 0 until numHashes) yield {
      val hashCodeRDD = colHashed1.map(ele => (coeffA(i) * ele + coeffB(i)) % nextPrime)
      hashCodeRDD.min.toInt // Track the lowest hash code seen.
    }

    val signature2 = for (i <- 0 until numHashes) yield {
      val hashCodeRDD = colHashed2.map(ele => (coeffA(i) * ele + coeffB(i)) % nextPrime)
      hashCodeRDD.min.toInt // Track the lowest hash code seen
    }

    val count = (0 until numHashes)
      .map(k => if (signature1(k) == signature2(k)) 1 else 0)
      .fold(0)(_ + _)


    val jSimilarity = count / numHashes.toDouble
    jSimilarity
  }


  //  def approach1(list1Values: List[String], list2Values: List[String]) = {
  //    val colHashed1 = list1Values.toSet
  //    val colHashed2 = list2Values.toSet
  //
  //    val jSimilarity = colHashed1.intersection(colHashed2).distinct.count / (colHashed1.union(colHashed2).distinct.count.toDouble)
  //    jSimilarity
  //  }


  def approach1(list1Values: List[String], list2Values: List[String]) = {
    val colHashed1 = list1Values.toSet
    val colHashed2 = list2Values.toSet

    val jSimilarity = (colHashed1 & colHashed2).size / (colHashed1 ++ colHashed2).size.toDouble
    jSimilarity
  }

  def main(args: Array[String]) {

    val list1Values = List("a", "b", "c")
    val list2Values = List("a", "b", "d")

    for (i <- 0 until 5) {
      println(s"Iteration ${i}")
      println(s" - Approach 1: ${approach1(list1Values, list2Values)}")
      println(s" - Approach 2: ${approach2(list1Values, list2Values)}")
    }

  }
}

输出:

Iteration 0
 - Approach 1: 0.5
 - Approach 2: 0.5
Iteration 1
 - Approach 1: 0.5
 - Approach 2: 0.5
Iteration 2
 - Approach 1: 0.5
 - Approach 2: 0.8
Iteration 3
 - Approach 1: 0.5
 - Approach 2: 0.8
Iteration 4
 - Approach 1: 0.5
 - Approach 2: 0.4

你为什么用它?

【讨论】:

  • 方法 1 给出了 jaccard 相似度的准确值。方法 2 是一个近似值。预计方法 2 的计算比方法 1 简单得多。如果我们在方法 2 中调整“numHashes”值,我们可以让近似值更接近我认为方法 1 给出的实际 jaccard 相似度。问题是近似计算似乎比精确计算密集
【解决方案2】:

在我看来,minHashing 方法的开销成本刚刚超过了它在 Spark 中的功能。特别是随着numHashes 的增加。 以下是我在您的代码中发现的一些观察结果:

首先,while (randList.contains(randIndex)) 这部分肯定会随着 numHashes(顺便说一下等于 randList 的大小)的增加而减慢您的进程。

其次,重写这段代码可以节省一些时间:

var signature1 = Array.fill(numHashes){0}
for (i <- 0 to numHashes-1)
{
    // Evaluate the hash function.
    val hashCodeRDD = colHashed1.map(ele => ((coeffA(i) * ele + coeffB(i)) % nextPrime))

    // Track the lowest hash code seen.
    signature1(i) = hashCodeRDD.min.toInt
}

var signature2 = Array.fill(numHashes){0}
for (i <- 0 to numHashes-1)
{
    // Evaluate the hash function.
    val hashCodeRDD = colHashed2.map(ele => ((coeffA(i) * ele + coeffB(i)) % nextPrime))

    // Track the lowest hash code seen.
    signature2(i) = hashCodeRDD.min.toInt
}


var count = 0
// Count the number of positions in the minhash signature which are equal.
for(k <- 0 to numHashes-1)
{
  if(signature1(k) == signature2(k))
    count = count + 1
}  

进入

var count = 0
for (i <- 0 to numHashes - 1)
{
    val hashCodeRDD1 = colHashed1.map(ele => ((coeffA(i) * ele + coeffB(i)) % nextPrime))
    val hashCodeRDD2 = colHashed2.map(ele => ((coeffA(i) * ele + coeffB(i)) % nextPrime))

    val sig1 = hashCodeRDD1.min.toInt
    val sig2 = hashCodeRDD2.min.toInt

    if (sig1 == sig2) { count = count + 1 }
}

此方法将三个循环简化为一个。但是,我不确定这是否会大大增加计算时间。

我的另一个建议是,假设第一种方法仍然更快,那就是使用集合的属性来修改第一种方法:

val colHashed1_dist = colHashed1.distinct
val colHashed2_dist = colHashed2.distinct
val intersect_cnt = colHashed1_dist.intersection(colHashed2_dist).distinct.count

val jSimilarity = intersect_cnt/(colHashed1_dist.count + colHashed2_dist.count - intersect_cnt).toDouble

这样,您可以重用交集的值,而不是获得联合。

【讨论】:

  • 我同意你的观点。但正如你所说,代码组织的变化并没有产生任何重大影响。我仍然想知道 spark 中是否有更好的 minHash 实现
【解决方案3】:

实际上,在 LSH 方法中,您只需为每个文档计算一次 minHash,然后为每个可能的文档对比较两个 minHases。如果采用简单的方法,您将对每对可能的文档进行完整的文档比较。这大约是 N^2/2 次比较。因此,对于足够多的文档,计算 minHashes 的额外成本可以忽略不计。

您实际上应该比较琐碎方法的性能:

val jSimilarity = colHashed1.intersection(colHashed2).distinct.count/(colHashed1.union(colHashed2).distinct.count.toDouble)

Jaccard 距离计算的性能(代码中的最后几行):

var count = 0
// Count the number of positions in the minhash signature which are equal.
for(k <- 0 to numHashes-1)
{
  if(signature1(k) == signature2(k))
    count = count + 1
}  
val jSimilarity = count/numHashes.toDouble

【讨论】:

    猜你喜欢
    • 2022-01-04
    • 2017-03-27
    • 2018-01-23
    • 1970-01-01
    • 1970-01-01
    • 2019-03-26
    • 2023-02-15
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多