【问题标题】:Achieve concurrency when saving to a partitioned parquet file保存到分区 parquet 文件时实现并发
【发布时间】:2023-03-20 10:25:01
【问题描述】:

当使用partitionBydataframe 写入parquet 时:

df.write.partitionBy("col1","col2","col3").parquet(path)

我希望每个正在写入的分区都由一个单独的任务独立完成,并且与分配给当前 spark 作业的工作人员数量并行。

但是,在写入 parquet 时,实际上一次只有一个工作人员/任务在运行。那个工作人员正在循环遍历每个分区并连续写出.parquet 文件。为什么会出现这种情况 - 有没有办法在这个 spark.write.parquet 操作中强制并发?

以下是不是我想看到的(应该是700%+..)

从其他帖子中,我也尝试在前面添加repartition

Spark parquet partitioning : Large number of files

df.repartition("col1","col2","col3").write.partitionBy("col1","col2","col3").parquet(path)

不幸的是,这没有效果:仍然只有一名工人..

注意:我使用local[8]local 模式下运行,并且看到其他 spark 操作使用多达 8 个并发工作人员并使用高达 750% 的 cpu。 p>

【问题讨论】:

    标签: scala apache-spark parquet


    【解决方案1】:

    简而言之,从单个任务写入多个输出文件不是并行化的,但假设您有多个任务(多个输入拆分),每个任务都会在工作线程上获得自己的核心。

    写出分区数据的目的不是并行化您的写操作。 Spark 已经通过同时写出多个任务来做到这一点。目标只是优化未来的读取操作,您只需要保存数据的一个分区。

    在 Spark 中写入分区的逻辑设计为在将前一个阶段的所有记录写入目的地时仅读取一次。我相信部分设计选择也是为了防止出现以下情况 一个分区键有很多很多值。

    编辑:Spark 2.x 方法

    在 Spark 2.x 中,它通过分区键对每个任务中的记录进行排序,然后遍历它们,一次写入一个文件句柄。我假设他们这样做是为了确保如果您的分区键中有很多不同的值,他们永远不会打开大量文件句柄。

    供参考,排序如下:

    https://github.com/apache/spark/blob/master/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/FileFormatWriter.scala#L121

    向下滚动一点,您会看到它调用write(iter.next()) 循环遍历每一行。

    这是实际的写入(一次一个文件/分区键):

    https://github.com/apache/spark/blob/master/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/FileFormatWriter.scala#L121

    你可以看到它一次只打开一个文件句柄。

    编辑:Spark 1.x 方法

    spark 1.x 对给定任务所做的是循环遍历所有记录,每当遇到此任务之前从未见过的新输出分区时打开一个新文件句柄。然后它立即将记录写入该文件句柄并进入下一个。这意味着在任何给定时间,在处理单个任务时,它最多可以为该任务打开 N 个文件句柄,其中 N 是输出分区的最大数量。为了更清楚起见,这里有一些 python 伪代码来展示总体思路:

    # To write out records in a single InputSplit/Task
    handles = {}
    for row in input_split:
        partition_path = determine_output_path(row, partition_keys)
        if partition_path not in handles:
            handles[partition_path] = open(partition_path, 'w')
    
        handles[partition_path].write(row)
    

    上述写出记录的策略有一个警告。在 spark 1.x 中,设置 spark.sql.sources.maxConcurrentWrites 对每个任务可以打开的掩码文件句柄设置了上限。在此之后,Spark 将改为按分区键对数据进行排序,因此它可以遍历记录,一次写出一个文件。

    【讨论】:

    • 我没有把它当作一个“策略”:例如hive 使用分区来驱动读写的并发性。这些是不同的目录,对并行写入没有任何障碍。我的意思是最后 - 如果是这种情况,那么我可以编写自己的例程,通过mapPartitions 将每个分区的写入并行化到不同的工作人员中。
    • 即使它们是不同的文件夹,仍然需要考虑并发文件句柄的数量,以及您可以执行的最大总写入 IO。理论上,您可以使用 mapPartitions 手动并行编写任务,但您也可以只增加执行器或核心的数量来实现相同的目标(假设您的工作有足够的任务)
    • 更新了 Spark 2 的更多具体信息
    • 这是很好的研究。将问题保持更长的时间,以期有人放弃解决方法。 afa # of workers:OP 提到他们是 8 人,但只有一个正在使用。对于所有这些使用单个线程被认为是可以接受的,我感到非常惊讶。
    • 只有一名工人正在写拼花地板。正如 OP 中所提到的,local[8]other 火花任务中运行良好(700+%):所以问题在于拼花书写。把这部作品连载是没有意义的。
    猜你喜欢
    • 2015-11-28
    • 2020-03-21
    • 1970-01-01
    • 1970-01-01
    • 2022-12-17
    • 1970-01-01
    • 2022-01-20
    • 2019-10-29
    • 2017-12-02
    相关资源
    最近更新 更多