【发布时间】: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