代码之家  ›  专栏  ›  技术社区  ›  Powers

用Spark的partitionBy方法对S3中一个大的倾斜数据集进行划分

  •  1
  • Powers  · 技术社区  · 7 年前

    partitionBy 我尝试过的两种方法中,算法都在挣扎。

    问题1 :

    当我使用重新分区之前 repartitionBy ,Spark将所有分区作为单个文件写入,即使是大型分区

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

    :

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

    这段代码会写出大量的文件。

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

    当我在一个生产数据集上运行这个程序时,一个有1.3gb数据的分区被写成3100个文件。

    我想要什么

    1 回复  |  直到 7 年前
        1
  •  7
  •   Nick Chammas    5 年前

    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")
      
        2
  •  2
  •   Pedro Jofre-Lora    5 年前

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

    目前接受的答案在大多数情况下可能已经足够好了,但并不能满足将0.3gb分区写入单个文件的请求。相反,它会写出来 numPartitions

    您要寻找的是一种根据数据分区的大小动态调整输出文件数量的方法。为此,我们将以10465355的使用方法为基础 rand() repartition() ,并缩放 兰德() 基于该分区所需的文件数。

    我将用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 id=2 获取1个分区,然后 id={3,4,5}

    这个解决方案平衡了输出文件的大小,而不考虑数据倾斜,并且不限制并行性 relying on maxRecordsPerFile

    推荐文章