【问题标题】:Comparing intersection between two nodes using broadcast variable and using RDD.filter in Spark GraphX在 Spark GraphX 中使用广播变量和 RDD.filter 比较两个节点之间的交集
【发布时间】:2020-03-02 17:07:26
【问题描述】:

我在 GraphX 中处理图表。通过使用下面的代码,我创建了一个变量来存储 RDD 中节点的邻居:

val all_neighbors: VertexRDD[Array[VertexId]] = graph.collectNeighborIds(EdgeDirection.Either)

我使用广播变量通过以下代码向所有从站广播邻居:

val broadcastVar = all_neighbors.collect().toMap
val nvalues = sc.broadcast(broadcastVar)

我想计算两个节点邻居之间的交集。例如节点 1 和节点 2 邻居之间的交集。

起初我使用这段代码来计算使用广播变量 nvalues 的交集:

val common_neighbors=nvalues.value(1).intersect(nvalues.value(2))

一旦我使用下面的代码来计算两个节点的交集:

val common_neighbors2=(all_neighbors.filter(x=>x._1==1)).intersection(all_neighbors.filter(x=>x._1==2))

我的问题是:上述哪种方法更高效、更分布式和并行?使用广播变量nvalue计算交集还是使用过滤RDD方法?

【问题讨论】:

    标签: scala apache-spark spark-graphx


    【解决方案1】:

    我认为这取决于情况。

    如果您的nvalues 大小较小并且可以适合每个执行程序和驱动程序节点,则广播方法将是最佳的,因为数据缓存在执行程序中并且不会一遍又一遍地重新计算这些数据。此外,它将节省 spark 巨大的通信和计算负担。在这种情况下,另一种方法不是最优的,因为可能会发生all_neighbours rdd 每次都计算,这会降低性能,因为会有大量的重新计算并会增加计算成本。

    如果您的nvalues 无法适应每个执行程序和驱动程序节点, 广播将不起作用,因为它会引发错误。因此,别无选择,只能使用第二种方法,尽管它仍然可能导致性能问题,至少代码可以工作!!

    如果有帮助请告诉我!!

    【讨论】:

    • 是的,它帮助了我。谢谢你。我认为 nvalues 可以适合每个奴隶的记忆。因为我正在处理大约 300 万个节点的图。我觉得合适。如果我使用 nvalues,我可以再次对 nvalues 结果进行分布式和并行操作,例如 foreach 吗?例如,如果我使用 nvalues 来计算两个节点的交集,那么我可以使用 common_neighbors.foreach 对公共邻居的每个元素进行分布式操作吗?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多