【问题标题】:Partitioning a large skewed dataset in S3 with Spark's partitionBy method使用 Spark 的 partitionBy 方法对 S3 中的大型倾斜数据集进行分区
【发布时间】:2019-04-01 21:22:20
【问题描述】:

我正在尝试使用 Spark 将大型分区数据集写入磁盘,而 partitionBy 算法在我尝试过的两种方法中都遇到了困难。

分区严重倾斜 - 一些分区很大,而另一些则很小。

问题 #1

当我在repartitionBy 之前使用repartition 时,Spark 会将所有分区写入单个文件,甚至是大文件

val df = spark.read.parquet("some_data_lake")
df
  .repartition('some_col).write.partitionBy("some_col")
  .parquet("partitioned_lake")

这需要永远执行,因为 Spark 不会并行写入大分区。如果其中一个分区有 1TB 的数据,Spark 会尝试将整个 1TB 的数据写入单个文件。

问题 #2

当我不使用repartition 时,Spark 会写出太多文件。

这段代码会写出数量惊人的文件。

df.write.partitionBy("some_col").parquet("partitioned_lake")

我在一个很小的 ​​8 GB 数据子集上运行了这个,Spark 写出了 85,000 多个文件!

当我尝试在生产数据集上运行它时,一个包含 1.3 GB 数据的分区被写成 3,100 个文件。

我想要什么

我希望每个分区都写成 1 GB 的文件。因此,具有 7 GB 数据的分区将作为 7 个文件写出,而具有 0.3 GB 数据的分区将作为单个文件写出。

我最好的前进道路是什么?

【问题讨论】:

标签: apache-spark apache-spark-sql partitioning


【解决方案1】:

最简单的解决方案是在repartition 中添加一列或多列,并显式设置分区数。

val numPartitions = ???

df.repartition(numPartitions, $"some_col", $"some_other_col")
 .write.partitionBy("some_col")
 .parquet("partitioned_lake")

地点:

  • numPartitions - 应该是写入分区目录的所需文件数的上限(实际数字可以更低)。
  • $"some_other_col"(和可选的附加列)应该具有高基数并且独立于$"some_column(这两者之间应该存在函数依赖关系,并且不应该高度相关)。

    如果数据不包含这样的列,您可以使用o.a.s.sql.functions.rand

    import org.apache.spark.sql.functions.rand
    
    df.repartition(numPartitions, $"some_col", rand)
      .write.partitionBy("some_col")
      .parquet("partitioned_lake")
    

【讨论】:

  • 我有完全相同的情况,所以我选择了盐渍,我希望每个分区不超过 5 个文件。但是让我感到困惑的是,我在保存 Stage 时看到了 5 个任务,而我不太明白的地方是我认为数据应该是平衡的,但它应该以最能利用集群资源的方式进行组织
  • @10465355 你能扩展一下“这两者之间应该存在功能依赖”吗?谢谢!
  • 鉴于 rand 是可以接受的替代品,我应该补充一下。
【解决方案2】:

我希望每个分区都写成 1 GB 的文件。因此,具有 7 GB 数据的分区将作为 7 个文件写出,而具有 0.3 GB 数据的分区将作为单个文件写出。

目前接受的答案在大多数情况下可能已经足够好,但并不能完全满足将 0.3 GB 分区写入单个文件的请求。相反,它将为每个输出分区目录写出numPartitions 文件,包括 0.3 GB 分区。

您正在寻找一种通过数据分区大小动态扩展输出文件数量的方法。为此,我们将在 10465355 的方法的基础上使用 rand() 来控制 repartition() 的行为,并根据我们希望该分区的文件数量来扩展 rand() 的范围。

很难通过输出文件大小来控制分区行为,因此我们将使用每个输出文件所需的大致行数来控制它。

我将在 Python 中提供一个演示,但方法在 Scala 中基本相同。

from pyspark.sql import SparkSession
from pyspark.sql.functions import rand

spark = SparkSession.builder.getOrCreate()
skewed_data = (
    spark.createDataFrame(
        [(1,)] * 100 + [(2,)] * 10 + [(3,), (4,), (5,)],
        schema=['id'],
    )
)
partition_by_columns = ['id']
desired_rows_per_output_file = 10

partition_count = skewed_data.groupBy(partition_by_columns).count()

partition_balanced_data = (
    skewed_data
    .join(partition_count, on=partition_by_columns)
    .withColumn(
        'repartition_seed',
        (
            rand() * partition_count['count'] / desired_rows_per_output_file
        ).cast('int')
    )
    .repartition(*partition_by_columns, 'repartition_seed')
)

这种方法将平衡输出文件的大小,无论分区大小有多么倾斜。每个数据分区都会获得它需要的文件数量,以便每个输出文件具有大致请求的行数。

这种方法的先决条件是计算每个分区的大小,您可以在partition_count 中看到。如果您真的想动态扩展每个分区的输出文件数量,这是不可避免的。

为了证明这是正确的,让我们检查分区内容:

from pyspark.sql.functions import spark_partition_id

(
    skewed_data
    .groupBy('id')
    .count()
    .orderBy('id')
    .show()
)

(
    partition_balanced_data
    .select(
        *partition_by_columns,
        spark_partition_id().alias('partition_id'),
    )
    .groupBy(*partition_by_columns, 'partition_id')
    .count()
    .orderBy(*partition_by_columns, 'partition_id')
    .show(30)
)

这是输出的样子:

+---+-----+
| id|count|
+---+-----+
|  1|  100|
|  2|   10|
|  3|    1|
|  4|    1|
|  5|    1|
+---+-----+

+---+------------+-----+
| id|partition_id|count|
+---+------------+-----+
|  1|           7|    9|
|  1|          49|    6|
|  1|          53|   14|
|  1|         117|   12|
|  1|         126|   10|
|  1|         136|   11|
|  1|         147|   15|
|  1|         161|    7|
|  1|         177|    7|
|  1|         181|    9|
|  2|          85|   10|
|  3|          76|    1|
|  4|         197|    1|
|  5|          10|    1|
+---+------------+-----+

根据需要,每个输出文件大约有 10 行。 id=1 得到 10 个分区,id=2 得到 1 个分区,id={3,4,5} 每个得到 1 个分区。

此解决方案平衡了输出文件的大小,不受数据倾斜的影响,并且不会限制relying on maxRecordsPerFile 的并行度。

【讨论】:

    【解决方案3】:

    Nick Chammas 方法的替代方法是创建一个按主分区键分区的 row_number() 列,然后将其除以您希望在每个分区中出现的确切记录数。用 SPARK SQL 表示如下:

    SELECT /*+ REPARTITION(id, file_num) */
      id,
      FLOOR(ROW_NUMBER() OVER(PARTITION BY id ORDER BY NULL) / rows_per_file) AS file_num
    FROM skewed_data
    
    

    这样做的额外好处是,它允许您通过在辅助键上使用ORDER BY 子句,将大部分数据集中在一个分区中的文件中。如果与辅助键关联的行号跨越两个file_num 值,则不能保证辅助键位于同一位置。也有可能,而且实际上有点可能,最终得到一个文件,每个分区中的记录很少。

    【讨论】:

    • 最大的好处不是托管,因为你不能保证边界出现在哪里。最大的好处是这种方法比 Nick Chammas 的方法使用的阶段少一个,因此性能更高。添加到性能增益的是缺少连接,这意味着更少的相等性检查。
    猜你喜欢
    • 2023-03-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-05-24
    • 1970-01-01
    • 2014-09-08
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多