【问题标题】:Count occurrences of each item in a Scala parallel collection计算 Scala 并行集合中每个项目的出现次数
【发布时间】:2013-08-22 16:51:18
【问题描述】:

我的问题与Count occurrences of each element in a List[List[T]] in Scala 非常相似,只是我希望有一个涉及parallel collections 的有效解决方案。

具体来说,我有一个大 (~10^7) 向量 vec 的短 (~10) 个 Int 列表,我想为每个 Int x 获取 x 出现的次数,例如Map[Int,Int]。不同整数的数量为 10^6。

由于需要在机器上完成这项工作,它具有相当数量的内存 (150GB) 和内核数 (>100),因此并行集合似乎是一个不错的选择。下面的代码是一个好方法吗?

val flatpvec = vec.par.flatten
val flatvec = flatpvec.seq
val unique = flatpvec.distinct
val counts = unique map (x => (x -> flatvec.count(_ == x)))
counts.toMap

或者有更好的解决方案吗?如果您对 .seq 转换感到疑惑:由于某种原因,以下代码似乎没有终止,即使对于小示例也是如此:

val flatpvec = vec.par.flatten
val unique = flatpvec.distinct
val counts = unique map (x => (x -> flatpvec.count(_ == x)))
counts.toMap

【问题讨论】:

    标签: scala collections parallel-processing


    【解决方案1】:

    这有什么作用。 aggregatefold 类似,只是您还合并了顺序折叠的结果。

    更新:.par.groupBy 中存在开销并不奇怪,但我对常数因素感到惊讶。根据这些数字,你永远不会这样计算。另外,我不得不提高内存。

    用于构建从the overview 链接的结果映射is described in this paper 的有趣技术。 (它巧妙地保存了中间结果,然后在最后将它们并行合并。)

    但是,如果您真正想要的只是一个计数,那么复制 groupBy 的中间结果会很昂贵。

    这些数字是比较顺序的groupBy、并行,最后是aggregate

    apm@mara:~/tmp$ scalacm countints.scala ; scalam -J-Xms8g -J-Xmx8g -J-Xss1m countints.Test
    GroupBy: Starting...
    Finished in 12695
    GroupBy: List((233,10078), (237,20041), (268,9939), (279,9958), (315,10141), (387,9917), (462,9937), (680,9932), (848,10139), (858,10000))
    Par GroupBy: Starting...
    Finished in 51481
    Par GroupBy: List((233,10078), (237,20041), (268,9939), (279,9958), (315,10141), (387,9917), (462,9937), (680,9932), (848,10139), (858,10000))
    Aggregate: Starting...
    Finished in 2672
    Aggregate: List((233,10078), (237,20041), (268,9939), (279,9958), (315,10141), (387,9917), (462,9937), (680,9932), (848,10139), (858,10000))
    

    测试代码中没有什么神奇之处。

    import collection.GenTraversableOnce
    import collection.concurrent.TrieMap
    import collection.mutable
    
    import concurrent.duration._
    
    trait Timed {
      def now = System.nanoTime
      def timed[A](op: =>A): A =  {
        val start = now
        val res = op
        val end = now
        val lapsed = (end - start).nanos.toMillis
        Console println s"Finished in $lapsed"
        res
      }
      def showtime(title: String, op: =>GenTraversableOnce[(Int,Int)]): Unit = {
        Console println s"$title: Starting..."
        val res = timed(op)
        //val showable = res.toIterator.min   //(res.toIterator take 10).toList
        val showable = res.toList.sorted take 10
        Console println s"$title: $showable"
      }
    }
    

    它会生成一些感兴趣的随机数据。

    object Test extends App with Timed {
    
      val upto = math.pow(10,6).toInt
      val ran = new java.util.Random
      val ten = (1 to 10).toList
      val maxSamples = 1000
      // samples of ten random numbers in the desired range
      val samples = (1 to maxSamples).toList map (_ => ten map (_ => ran nextInt upto))
      // pick a sample at random
      def anyten = samples(ran nextInt maxSamples)
      def mag = 7
      val data: Vector[List[Int]] = Vector.fill(math.pow(10,mag).toInt)(anyten)
    

    aggregate的顺序操作和组合操作从一个任务中调用,结果赋值给一个volatile var。

      def z: mutable.Map[Int,Int] = mutable.Map.empty[Int,Int]
      def so(m: mutable.Map[Int,Int], is: List[Int]) = {
        for (i <- is) {
          val v = m.getOrElse(i, 0)
          m(i) = v + 1
        }
        m
      }
      def co(m: mutable.Map[Int,Int], n: mutable.Map[Int,Int]) = {
        for ((i, count) <- n) {
          val v = m.getOrElse(i, 0)
          m(i) = v + count
        }
        m
      }
      showtime("GroupBy", data.flatten groupBy identity map { case (k, vs) => (k, vs.size) })
      showtime("Par GroupBy", data.flatten.par groupBy identity map { case (k, vs) => (k, vs.size) })
      showtime("Aggregate", data.par.aggregate(z)(so, co))
    }
    

    【讨论】:

    • 有趣,但这不会导致为data 中的每个元素创建一个地图吗?
    • @mitchus z 可变更有意义,因此每个顺序操作一个映射,这是一个单线程任务,但我懒得解决它。我会把它放在我的任务中。
    • @mitchus 更新为使用可变结果,这很有效。看到惊人的数字。或者也许他们并不奇怪。
    • 我还不知道聚合。似乎是比我建议的更好的解决方案。
    • 确实,看起来很不错。
    【解决方案2】:

    如果您想使用并行集合和 Scala 标准工具,您可以这样做。按标识对您的集合进行分组,然后将其映射到 (Value, Count):

    scala> val longList = List(1, 5, 2, 3, 7, 4, 2, 3, 7, 3, 2, 1, 7)
    longList: List[Int] = List(1, 5, 2, 3, 7, 4, 2, 3, 7, 3, 2, 1, 7)                                                                                            
    
    scala> longList.par.groupBy(x => x)
    res0: scala.collection.parallel.immutable.ParMap[Int,scala.collection.parallel.immutable.ParSeq[Int]] = ParMap(5 -> ParVector(5), 1 -> ParVector(1, 1), 2 -> ParVector(2, 2, 2), 7 -> ParVector(7, 7, 7), 3 -> ParVector(3, 3, 3), 4 -> ParVector(4))                                                                     
    
    scala> longList.par.groupBy(x => x).map(x => (x._1, x._2.size))
    res1: scala.collection.parallel.immutable.ParMap[Int,Int] = ParMap(5 -> 1, 1 -> 2, 2 -> 3, 7 -> 3, 3 -> 3, 4 -> 1)                                           
    

    甚至更好,如 cmets 中建议的 pagoda_5b:

    scala> longList.par.groupBy(identity).mapValues(_.size)
    res1: scala.collection.parallel.ParMap[Int,Int] = ParMap(5 -> 1, 1 -> 2, 2 -> 3, 7 -> 3, 3 -> 3, 4 -> 1)
    

    【讨论】:

    • 好主意,我会试试的。
    • 您可以通过使用identity(_) 函数作为groupBy 的参数而不是x =&gt; x 来进行一些小的改进。您也可以 map 使用 mapValues(_.size) 分组 Map 的值
    • 好建议,pagoda_5b。我将其添加到答案中。 :)
    • 只是检查一下,你知道mapValues 是一个非严格视图,所以size 是按需计算的。 scalapuzzlers.com/#pzzlr-037
    猜你喜欢
    • 1970-01-01
    • 2011-12-27
    • 1970-01-01
    • 2016-01-24
    • 1970-01-01
    • 2014-08-22
    • 1970-01-01
    • 2011-07-02
    • 2012-07-23
    相关资源
    最近更新 更多