【问题标题】:How to specify file size using repartition() in spark如何在 spark 中使用 repartition() 指定文件大小
【发布时间】:2021-04-30 22:03:59
【问题描述】:

我正在使用 pyspark,并且我有一个大型数据源,我想重新分区,明确指定每个分区的文件大小。

我知道使用repartition(500) 函数会将我的镶木地板分成大小几乎相等的 500 个文件。 问题是每天都会有新数据添加到这个数据源中。在某些日子可能会有很大的输入,而在某些日子可能会有较小的输入。因此,在查看一段时间内的分区文件大小分布时,每个文件的大小分布在 200KB700KB 之间。

我正在考虑指定每个分区的最大大小,这样无论文件数量如何,我每天每个文件的文件大小都差不多。 这将有助于我稍后在这个大型数据集上运行我的工作,以避免扭曲的执行器时间和洗牌时间等。

有没有办法使用repartition() 函数或在将数据帧写入镶木地板时指定它?

【问题讨论】:

    标签: apache-spark pyspark parquet partitioning


    【解决方案1】:

    您可以考虑使用参数maxRecordsPerFile 编写结果。

    storage_location = //...
    estimated_records_with_desired_size = 2000
    result_df.write.option(
         "maxRecordsPerFile", 
         estimated_records_with_desired_size) \
         .parquet(storage_location, compression="snappy")
    

    【讨论】:

    • 但要做到这一点,我首先需要找出一个 100MB 文件中有多少条记录,then set the maxRecordsPerFile 对吗?有没有办法直接指定文件的最大大小?
    • 你的理解是正确的。直接回答您的问题,不,目前不。内存中的 DataFrame 在写入磁盘(或 AWS S3 等对象存储位置)之前需要进行编码和压缩,默认持久模式为StorageLevel.MEMORY_AND_DISK。简单来说,在将结果完全写入磁盘之前,无法估计写入过程中文件的实际大小。
    • 明白。如果是这样,我将如何查找仅 100MB 数据的行数?
    • 假设您的连续数据结构是一致的,并且您有一个包含 1,000 条记录的文件(结果)。使用前提条件,您可以获得结果的平均行大小。假设平均大小为 100kb,那么 100 MB 的估计行数将为 (100 x 1,024) / 100 = 1024(行数)。压缩与否,.csv.parquet,进度差不多。这是满足您的要求或愿望的最快方式。在压缩场景 (file_name.csv.gz) 中,可能会有一个数学公式,但它会比我提供的上一个方法花费更多时间。
    • 谢谢。我在问,以便在我尝试在答案中使用您的方法时估计将进入分区的行数。你如何找到记录的大小?我尝试使用此论坛的答案之一中建议的getByte(df.head()),但它对 pyspark 无效。
    猜你喜欢
    • 2018-01-04
    • 2022-10-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-10-09
    • 1970-01-01
    • 2019-01-15
    相关资源
    最近更新 更多