【问题标题】:scala: use parallel collections to do a foreach, and then do something else on each partition?scala:使用并行集合做一个foreach,然后在每个分区上做其他事情?
【发布时间】:2011-08-26 18:37:00
【问题描述】:

我有一个Seq 的项目,我需要对每个项目做一些事情,然后执行不需要任何输入的最后一步。我想使用 par 来加快速度:将Seq 拆分为多个分区,并在每个分区内对每个项目执行某些操作并为每个分区执行最后一步。我希望在处理特定分区的线程中运行最后一步。有没有办法做到这一点? aggregate() 似乎做得不对。

这里有一些示例代码:

// non-parallel case
val mySeq = Seq[Item]
mySeq foreach { actOnItem(_) }
doFinalStep()

// ideal par case
val mySeq = Seq[Item]
mySeq.par foreachThenDoFinalStepAfterPartition { actOnItem(_), doFinalStep }

【问题讨论】:

  • 你能再描述一下吗?为什么每一批元素都处理完之后,还要做一些特别的事情?
  • @amir 该链接与我的用例无关。
  • @axel22 处理每个元素都会累积一些数据。我想在最后对这些数据做点什么。我可以轻松地移动它,以便它累积返回值,但我只想为每个分区处理一次累积。
  • 这听起来接近于 aggregate() 的作用......聚合的哪一部分没有达到你想要的效果?

标签: scala


【解决方案1】:

我假设actOnItem 会产生doFinalStep 使用的某种副作用。由于您同时有多个副作用组,现在您希望使用doFinalStep 处理每个组,因此您必须在一个并发数据结构中跟踪不同组的不同副作用。在不确切知道actOnItemdoFinalStep 的作用的情况下,很难将其转换为更实用的样式。你可以这样做:

class SF { /* whatever your sideeffect is */ }

val ac = new java.util.concurrent.atomic.AtomicInteger(0)
val sfmap = new java.util.concurrent.ConcurrentHashMap[Int, SF]()
def newSideeffectIndex() = {
  val i = ac.fetchAndIncrement()
  sfmap.put(i, new SF())
  i
}

val mySeq = Seq[Item]()
mySeq.aggregate(-1)((u, x) => actOnItem(u, x), (u1, u2) => {
  doFinalStep(u1)
  doFinalStep(u2)
})

def actOnItem(u0: Int, x: Item) {
  val u = if (u0 == -1) newSideeffectIndex() else u0

  // do whatever you need to do with `x`
  // ...

  val sf = map.get(u)
  // do something with `sf` - update it somehow based on `x`

  u
}

def doFinalStep(u: Int) {
  val sf = map.get(u)

  if (sf != null) {
    // do the final step here using `sf`
  }

  map.remove(u)
}

说明:每个新分区以聚合值 -1 开始。如果检测到 -1,聚合部分(第一个闭包)将初始化当前分区的聚合值。初始化将选择一个唯一的整数并创建一个包含副作用的 SF 对象。之后,它将处理当前元素,然后更新副作用。在缩减步骤中,您知道该分区已被处理,因此您可以对该分区执行最后一步 - 并将其从并发映射中删除。

【讨论】:

    猜你喜欢
    • 2015-11-01
    • 1970-01-01
    • 1970-01-01
    • 2021-08-29
    • 2015-10-20
    • 2014-09-14
    • 2014-12-03
    • 2021-12-31
    • 1970-01-01
    相关资源
    最近更新 更多