【问题标题】:joins and cogroup in Spark在 Spark 中加入和共同组
【发布时间】:2015-04-15 17:52:25
【问题描述】:

有迹象表明 Spark 中的连接是使用 / 基于 cogroup 函数/primitive/transform 实现的。所以让我首先关注 cogroup - 它返回一个 RDD 的结果,它基本上由 cogrouped RDD 的所有元素组成。换一种说法——对于每个同组 RDD 中的每个键,至少有一个来自同组 RDD 中的至少一个的元素。

这意味着当更小,而且流式传输,例如JavaPairDstreamRDDs 不断加入更大的批处理 RDD,这将导致为结果(cogrouped)RDD 的多个实例分配 RAM,a.k.a 本质上是大批处理 RDD 等等...... 显然,当 DStream RDD 被丢弃并且它们定期这样做时,RAM 将被返回,但这似乎仍然是 RAM 消耗中不必要的峰值

我有两个问题:

  1. 是否有更“精确”控制 cogroup 过程的方法,例如告诉它只包含共组 RDD 元素,其中每个给定键的共组 RDD 的每个元素中至少有一个元素。根据当前的 cogroup API,这是不可能的

  2. 1234563仍然是同样严重的 RAM 消耗

【问题讨论】:

    标签: apache-spark spark-streaming


    【解决方案1】:

    没那么糟。这在很大程度上取决于分区的粒度。 Cogroup 将首先在磁盘中按密钥洗牌到不同的执行程序节点。对于每个键,是的,对于两个 RDD,具有该键的所有元素的整个集合都将加载到 RAM 中并提供给您。但并非所有密钥都需要在任何给定时间都在 RAM 中,因此除非您的数据确实存在偏差,否则您不会因此受到太大影响。

    【讨论】:

    • 在协同分组帮助之前是否会使用相同的分区器重新分区?
    • 我有超过 5 个 JavaPairRDD,包含一个主 pairRDD。我想结合这些所有基于主pairRDD。我该怎么做?
    • 如何将cogroup 用于大型数据集,例如当我使用collect() 时,它会抛出内存不足异常rdd1 = rdd2.cogroup(rdd3).collect。你能帮忙解决这个问题吗[stackoverflow.com/questions/47180307/….可以分区帮助我是新手来激发任何帮助来解决这个问题。
    • @Vignesh 当然,我离开并在那边回答。
    猜你喜欢
    • 1970-01-01
    • 2018-01-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-04-08
    • 2017-04-24
    • 2015-05-05
    • 2016-01-24
    相关资源
    最近更新 更多