【问题标题】:Updating an Array via Scala Parallel Collections通过 Scala 并行集合更新数组
【发布时间】:2014-07-11 20:51:55
【问题描述】:

我将这个 HashMap 数组定义如下

var distinctElementsDefinitionMap: scala.collection.mutable.ArrayBuffer[HashMap[String, Int]] = new scala.collection.mutable.ArrayBuffer[HashMap[String, Int]](300) with scala.collection.mutable.SynchronizedBuffer[HashMap[String, Int]]

现在,我有 300 个元素的并行集合

val max_length = 300
val columnArray = (0 until max_length).toParArray
import scala.collection.parallel.ForkJoinTaskSupport
columnArray.tasksupport = new ForkJoinTaskSupport(new scala.concurrent.forkjoin.ForkJoinPool(100))
columnArray foreach(i => {
    // Do Some Computation and get a HashMap
    var distinctElementsMap: HashMap[String, Int] = //Some Value
    //This line might result in Concurrent Access Exception
    distinctElementsDefinitionMap.update(i, distinctElementsMap)
})

我现在正在上面定义的columnArray 上的foreach 循环内运行计算密集型任务。 计算完成后,我希望每个线程更新distinctElementsDefinitionMap 数组的特定条目。 每个线程只会更新特定的索引值,对执行它的线程来说是唯一的。 我想知道这个数组条目的更新是否安全,多个线程可能同时写入它? 如果没有,是否有 synchronized 这样做的方式,所以它是线程安全的? 谢谢!

更新: 看来这确实不是安全的方法。我收到了java.util.ConcurrentModificationException 在使用并行集合时如何避免这种情况的任何提示。

【问题讨论】:

  • 你在滥用并行集合——它并不意味着是一个时尚的普通线程池,而是你将处理移交给一个聪明的池(工作窃取 ftw!)并避免使用副作用和然后使用处理结果(可能以单线程方式)。再一次,它是一个 parallel 集合,而不是 concurrent。或许您可以向我们提供有关您要归档的内容的更大图景?
  • 我完全同意,我知道我正在做的不是最优化的方式,甚至不是好的方式。但我只是 Scala 的初学者,并且仍在寻找解决方法。但我需要一个并行循环,这是我想到的唯一方法。为这种基本方法道歉!
  • 不用担心,但不清楚为什么在完成每项任务后需要更新地图中的条目。如果你澄清它,也许我们可以想出一个替代的惯用解决方案。
  • 嗯,基本上我有一个 300 列,800 万行的数据集。我需要为每一列创建一个哈希图,为这个哈希图的每个不同值找到从字符串值到整数值的映射。因此需要一个 HashMap 数组。数组的每个条目都是对应于该列的不同值的哈希图。一种方法是按顺序查找每列的哈希图并更新distinctElementsDefinitionMap 数组。但我想加快速度,从而使用并行集合。
  • 编辑了我的问题以显示更新完成

标签: multithreading scala synchronization scala-collections parallel-collections


【解决方案1】:

使用.groupBy操作,据我判断it is parallelized(不像其他一些方法,比如.sorted

case class Row(a: String, b: String, c: String)
val data = Vector(
  Row("foo", "", ""), 
  Row("bar", "", ""), 
  Row("foo", "", "")
)

data.par.groupBy(x => x.a).seq
// Map(bar -> ParVector(Row(bar,,)), foo -> ParVector(Row(foo,,), Row(foo,,)))

希望你明白了。

或者,如果您的 RAM 允许您在每列而不是行上并行处理,它必须比您当前的方法更有效(更少争用)。

val columnsCount = 3 // 300 in your case
Vector.range(0, columnsCount).par.map { column => 
  data.groupBy(row => row(column))
}.seq 

尽管即使使用单列也可能会出现内存问题(8M 行可能很多)。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-11-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多