【问题标题】:How to use Broadcast variable correctly in Spark GraphX?如何在 Spark GraphX 中正确使用广播变量?
【发布时间】:2020-03-02 15:34:45
【问题描述】:

我使用 GraphX 来处理图形。我已经使用 GraphLoader 来加载它,并使用以下代码创建了一个包含每个节点的邻居的变量:

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

因为我经常需要节点邻居,所以我决定广播它们。当我使用此代码时出现错误:

val broadcastVar = sc.broadcast(all_neighbors)

但是当我使用这段代码时没有错误:

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

使用collect()进行广播是否正确?

还有一个问题。我想将此广播变量更改为键值。这段代码对吗?

val nvalues = broadcastVar.value.toMap

上面的代码(我的意思是nvalues)是否广播到集群中的所有从站?我也应该广播 nvalues 吗?我对广播主题有点困惑。请帮我解决这个问题。

【问题讨论】:

    标签: scala apache-spark spark-graphx


    【解决方案1】:

    有两个问题:

    使用collect()进行广播是否正确?

    all_neighbors 是 VertexRDD 类型,本质上是一个 RDD。您可以广播的 RDD 中没有任何内容。 RDD 是一种数据结构,描述了对某些数据集的分布式计算。通过 RDD 的特性,您可以描述计算什么以及如何计算。它是一个抽象实体。你只能广播一个真实的值,但 RDD 只是一个值的容器,只有在 executor 处理它们的数据时才可用。

    引用Broadcast Variables:

    广播变量允许程序员保留一个只读变量 缓存在每台机器上,而不是随任务一起发送一份副本。 例如,它们可以用来给每个节点一个大的副本 以高效的方式输入数据集。

    这意味着显式创建广播变量仅有用 当跨多个阶段的任务需要相同的数据或缓存时 反序列化形式的数据很重要。

    这就是我们需要执行collect RDD 持有的数据集的原因,它将 RDD 转换为本地可用的集合,然后可以广播。

    注意:当你执行collect操作时,数据会在驱动节点中累积,然后广播出去。所以如果驱动节点中的空间较小,就会抛出错误

    上面的代码(i 表示 nvalues)是否向所有从站广播 簇??我也应该广播 nvalues 吗?

    这完全取决于您的用例。如果您只想使用broadcastVar,则仅广播它,或者如果您想使用nvalues,则仅广播nvalues,否则您可以同时广播这两个值,但您需要注意内存限制。

    如果有帮助请告诉我!!

    【讨论】:

    • 非常感谢您的完整描述。这对我很有帮助。谢谢。
    • 如果有帮助,请投票并关闭此答案。非常感谢。
    • 我有另一个问题与我问的这个问题有相同的主题。你能看看那个问题吗?这是那个链接。 [链接] (stackoverflow.com/questions/60493554/…)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-08-18
    • 1970-01-01
    相关资源
    最近更新 更多