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