【问题标题】:Spark: Distributed removal/addition of elements in a set?Spark:集合中元素的分布式删除/添加?
【发布时间】:2016-06-01 16:09:35
【问题描述】:

我正在尝试将 ML 算法转换为 Spark Scala 以利用集群的强大功能。伪代码的相关位如下:

initialize set of elements

while(set not empty) {

  while(...) { remove a given element from the set }

  while(...) { add a given element to the set }

}

有没有办法并行化这样的事情?

我会直观地说这不能以分布式方式实现(迭代次数未知),但我一直在读到 Spark 允许实现迭代 ML 算法。

这是我目前尝试过的:

  • 最初使用可变 Set 并在简单 Scala 的循环期间删除/添加元素。它运行正常,但我觉得整个代码只会在驱动程序上执行,这限制了使用 Spark 的兴趣?
  • 将集合设置为 RDD,并在每次迭代期间将 var 替换为带有减去/添加元素的新 RDD(我认为这非常重?)。没有出现错误,但变量实际上并没有得到更新。

    mySetRDD = mySetRDD.subtract(sc.parallelize(Seq(element)))

  • 查找了 Accumulators 以在多个执行程序之间保持一组元素更新其内容(存在/不存在元素)的方法,但除了简单的数值更新之外,它们似乎不允许其他事情。

【问题讨论】:

  • 一般来说,您可以丢弃在转换中使用时不能提供强一致性保证的累加器。其余的取决于手头的特定问题。迭代性质在这里不是问题。
  • 我明白了。把累加器放在一边,我的集合可以实现什么样的实现?我可能是错的,但我觉得如果我保留我的可变集,整个执行将在一台机器上保持完全顺序,而不使用 Spark 的功能。另一方面,我尝试的 RDD“更新”似乎不起作用。
  • 这里真的没有通用的答案。如果您无法以不同于将任意元素从一个集合移动到另一个集合的方式来表达问题,那么就没有充分的理由使用 Spark。我的意思是你可以使用 IndexedRDD 或类似的工具,但这没有任何价值。您遇到的问题是算法设计本身没有编程。

标签: scala apache-spark


【解决方案1】:

创建 PairRDD 然后 repartitionByKey 说 x 个分区。 之后就可以使用了

PairRdd1.zipPartition() 获取 rdds 分区的迭代器。然后您可以编写一个函数,该函数将在两个迭代器上运行以生成第三个或输出迭代器。

由于您已按密钥对 rdd 进行了重新分区,因此您无需跟踪跨分区的删除。

https://spark.apache.org/docs/1.0.2/api/java/org/apache/spark/rdd/RDD.html#zipPartitions(org.apache.spark.rdd.RDD, boolean, scala.Function2, scala.reflect.ClassTag, scala.reflect.ClassTag)

【讨论】:

  • “你不需要跨分区跟踪删除”我不确定我明白了。
  • 我的意思是每个分区都有不同的元素,所以你不必担心其他分区中的数据。
  • 哦,对了!但是,如果需要删除的元素是整个集合中的最大值,那么真的没有办法通过分区来完成这项工作,是吗?
  • 你知道要移除的元素的值吗?或者你想删除最大值?
  • 具体的最大值。
猜你喜欢
  • 2015-08-13
  • 2022-10-23
  • 2011-02-06
  • 1970-01-01
  • 1970-01-01
  • 2011-05-30
  • 2021-07-07
  • 2022-12-18
  • 1970-01-01
相关资源
最近更新 更多