【问题标题】:Managing Spark partitions after DataFrame unions在 DataFrame 联合之后管理 Spark 分区
【发布时间】:2017-01-15 19:27:21
【问题描述】:

我有一个 Spark 应用程序,需要大量使用 unions,因此我将在不同时间、不同情况下将大量 DataFrame 联合在一起。我正在尝试尽可能高效地运行它。我对 Spark 还很陌生,我突然想到了一些事情:

如果我有具有 X 个分区 (numAPartitions) 的 DataFrame 'A' (dfA),然后我将它合并到具有 Y 个分区 (@987654325) 的 DataFrame 'B' (dfB) @),那么得到的联合 DataFrame (unionedDF) 会是什么样子,结果是分区?

// How many partitions will unionedDF have?
// X * Y ?
// Something else?
val unionedDF : DataFrame = dfA.unionAll(dfB)

对我来说,理解这一点似乎非常重要,因为 Spark 的性能似乎在很大程度上依赖于 DataFrame 采用的分区策略。因此,如果我要左右合并 DataFrame,我需要确保我不断管理合并后的 DataFrame 的分区。

唯一我能想到的事情(以便正确管理联合数据帧的分区)是重新分区它们,然后在我联合它们后立即将数据帧持久保存到内存/磁盘:

val unionedDF : DataFrame = dfA.unionAll(dfB)
unionedDF.repartition(optimalNumberOfPartitions).persist(StorageLevel.MEMORY_AND_DISK)

这样,一旦它们被联合起来,我们就会对它们进行重新分区,以便将它们正确地分布在可用的 worker/executor 上,然后 persist(...) 调用告诉 Spark 不要从内存中驱逐 DataFrame,所以我们可以继续努力。

问题是,重新分区听起来很昂贵,但它可能没有替代方法那么昂贵(根本不管理分区)。是否有关于如何在 Spark-land 中有效管理工会的普遍接受的指导方针?

【问题讨论】:

    标签: apache-spark distributed-computing partitioning spark-dataframe unions


    【解决方案1】:

    是的,分区对于 很重要。

    我想知道您是否可以通过拨打以下电话自己找出答案:

    yourResultedRDD.getNumPartitions()
    

    我必须坚持,后工会吗?

    一般来说,如果你要多次使用它,你必须持久化/缓存一个 RDD(不管它是联合的结果,还是土豆的结果 :))。这样做会阻止 在内存中再次获取它,并且在某些情况下可以将应用程序的性能提高 15%!

    例如,如果您打算只使用一次生成的 RDD,那么不持久化它是安全的。

    我必须重新分区吗?

    由于您不关心查找分区数,您可以阅读我的memoryOverhead issue in Spark ,了解分区数如何影响您的应用程序。

    一般来说,你拥有的分区越多,每个执行器处理的数据块就越小。

    回想一下,一个 worker 可以托管多个 executor,您可以将它想象成 worker 是集群的机器/节点,而 executor 是一个在该 worker 上运行的进程(在核心中执行)。

    Dataframe 不是一直在内存中吗?

    不是真的。这对于 来说非常可爱,因为当您处理 时,您不希望不必要的东西存在于内存中,因为这会威胁到您的应用程序的安全。

    DataFrame 可以存储在 为您创建的临时文件中,并且仅在需要时才加载到应用程序的内存中。

    更多阅读:Should I always cache my RDD's and DataFrames?

    【讨论】:

    • 谢谢@gsamaras (+1) - 我并不是真的在问如何找出我的 DataFrames 有多少个分区,我真正要问的是我是否必须重新分区并坚持这样至于管理他们,工会后。有什么想法吗?
    • @smeeb 我用一些想法更新了我的答案,因为你没有任何其他答案..希望有帮助! :)
    • 再次感谢@gsamaras (+1) - 如果您不介意的话,请回答两个快速的后续问题:(1) 当您说“这样做可以防止火花从内存中再次获取它并可以提高应用程序的性能...”,我想我不明白您所说的“...防止火花获取它in再次记忆​​>...”。 DataFrame 不是 always 在内存中吗?如果 Spark 正在 内存中获取 DataFrame,它还能将 DataFrame 加载到哪里?内存不是存储数据帧最快的地方吗?!?
    • (2) 我如何决定每个worker应该获得多少个executor,以及每个executor应该获得多少个partition?再次感谢!
    • @smeeb 我更新了我的答案。至于(2),该死的太宽泛了......如果你想问一个新问题......;)
    【解决方案2】:

    Union 只是将数据框 1 和数据框 2 中的分区数相加。两个数据框具有相同的列数和相同的顺序来执行联合操作。所以不用担心,如果两个数据帧中的分区列不同,则最多会有 m + n 个分区。

    你不需要在加入后重新分区你的数据帧,我的建议是使用合并代替重新分区,合并普通分区或合并一些小分区并避免/减少分区内的混洗数据。

    如果您在每次联合后缓存/持久化数据帧,您将降低性能,并且沿袭不会被缓存/持久性破坏,在这种情况下,垃圾收集将在一些大量内存密集型操作的情况下清理缓存/内存,并且重新计算会增加相同的计算时间,这一次可能是清除/删除数据需要部分计算。

    由于火花转换是惰性的,即; unionAll 是惰性操作,coalesce/repartition 也是惰性操作,并在第一次操作时起作用,因此请尝试在计数器为 8 的间隔后合并 unionall 结果并减少结果数据帧中的分区。如果您的解决方案中有大量内存密集型操作,请使用检查点来中断沿袭和存储数据。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-03-06
      • 2022-11-12
      • 2018-01-19
      • 2017-01-15
      • 1970-01-01
      • 1970-01-01
      • 2019-03-02
      相关资源
      最近更新 更多