【问题标题】:Apache Spark Shared CounterApache Spark 共享计数器
【发布时间】:2014-08-08 19:38:27
【问题描述】:

我是 Apache Spark 的新手,我正在使用 GraphX。所以我必须使用 Scala,我也是新手 ;-)。

更新

我有一个图表,比如下图:

每个节点都有自己的 HashMap 或 List 可以存储 ID。现在我正在遍历图的三元组,如果边属性匹配一个条件(在本例中被忽略),那么我想在这条边的开始和结束节点中存储相同的 ID。

该算法经过一轮后,结果可能如下所示:

这里是代码(缩短):

val newNodes = graph.triplets.flatMap(triplet => {
    val newId = Counter.getId();
    val map = List((srcId, newId), (dstId, newId))
    // Outputs sth. like
    //  (1, 1)   (3, 1)   (2, 2)    (3, 2)
}

我从计数器对象中获取唯一 ID:

object Counter{
    private var resultCount: Integer = 0;

    def getResultID(): Integer = {
        resultCount = resultCount + 1;
        return resultCount;
    }
}

在 flatMap 之后,我按节点 id 对所有元组进行分组,然后将一个节点的所有 id 放入列表中(使用地图运算符)。所以结果是节点 3:(3, List(1, 2))。然后,此结果会通过 outerJoin 存储回图表。

所以我的问题是,我是否必须关心,通过同步方法,ID 是唯一的,或者以这种方式可以吗?如果有人有另一个想法通过解决整个问题而不给出明确的 ID,例如使用 zip 方法,那么这也很好:-)。

除了这个问题,有人可以解释一下,Counter 对象在运行时发生了什么吗?因为它是一个单例,它是否驻留在执行驱动程序的某个地方(在 Master 上?),因为我在某处读过,您可以在正常代码中使用然后在使用 Spark 进行并行计算时使用的变量,被复制到Worker/Thread,这里不应该发生。

提前致谢!

【问题讨论】:

  • 如果您在计算 RDD 中的内容,您应该使用 accumulators
  • 谢谢,但我不想计算 RDD 中的内容。相反,我想在使用 RDD(发生分布式)做某事时给事物一个唯一的 ID。所以我不知道我是否必须确保 ID 是唯一的,并且 2 个或更多线程可以(可能)获得相同的一个(丢失更新问题等等)。
  • 啊,这有一个简单的解决方案,给我一分钟发帖

标签: scala apache-spark


【解决方案1】:

没有必要自己实现这个,分配唯一的 id 是如此普遍,以至于 spark 已经在 def zipWithUniqueId(): RDD[(T, Long)] 中内置了它,你可以看到它为每个值分配一个唯一的 long,这意味着它返回一个元组的 RDD .示例用法:

val uniqIds = vertexData.zipWithUniqueId().map((k,v)=>(v,k)) //I'm assuming you want the unique ids as the vertexId

你也可以使用边缘属性来做到这一点

【讨论】:

  • 谢谢。在考虑了 zip 的解决方案后,它对我有用。
  • @th0rsch 很高兴我的回答对我忽略了您的其他评论感到抱歉
猜你喜欢
  • 1970-01-01
  • 2011-10-11
  • 1970-01-01
  • 1970-01-01
  • 2017-04-24
  • 2011-04-30
  • 1970-01-01
  • 2010-12-27
  • 2011-01-06
相关资源
最近更新 更多