【问题标题】:How to consolidate results of a spark SQL query to avoid lots of small files / avoid empty files如何合并 Spark SQL 查询的结果以避免大量小文件/避免空文件
【发布时间】:2018-04-06 12:58:30
【问题描述】:

上下文:在我们的数据管道中,我们使用 spark SQL 运行大量查询,这些查询由我们的最终用户作为文本文件提供,然后我们对其进行参数化。

情况

我们的查询如下所示:

INSERT OVERWRITE TABLE ... PARTITION (...)

SELECT 
  stuff
FROM
  sometable

问题是,当您查看此结果时,它不会创建一堆大小为最大块大小的文件,而是创建 200 个小文件(因为默认情况下 spark 创建 200 个分区)。 (对于某些查询,取决于输入数据和SELECT 查询,对于 200 百阅读数不胜数)。大量的小文件让我们不受系统管理员的欢迎。

已尝试修复(不起作用)

大量文档表明,在这种情况下,您应该使用 DISTRIBUTE BY 以确保给定分区的所有数据都进入同一个分区,所以让我们尝试以下操作:

INSERT OVERWRITE TABLE ... PARTITION (...)

SELECT 
  stuff
FROM
  sometable
DISTRIBUTE BY
  1

那么为什么这不起作用(在 spark 2.0 和 spark 2.2 上测试)?它确实成功地将所有数据发送到一个reducer - 所有实际数据都在一个大文件中。但它仍然会创建 200 个文件,其中 199 个是空的! (我知道我们可能应该 DISTRIBUTE BY 我们的分区列,但这是提供最简单的示例)

修复确实有效,但不适合我们的用例

可以通过使用coalescepartition 来做正确的事情,因此(在pyspark 语法中):

select = sqlContext.sql('''SELECT stuff FROM sometable''').coalesce(1)
select.write.insertInto(target_table, overwrite=True)

但我不想这样做,因为我们需要彻底改变用户向我们提供查询的方式。

我还看到我们可以设置:

conf.set("spark.sql.shuffle.partitions","1");

但我还没有尝试过,因为我不想强制(相当复杂的)查询中的所有计算都发生在一个 reducer 上,只发生在最后写入磁盘的那个上。 (如果我不应该担心这个,请告诉我!)

问题

  • 仅使用 spark SQL 语法,我如何编写一个尽可能少写入文件且不创建大量空/小文件的查询?

可能相关:

【问题讨论】:

  • 没有回答您的问题,但我们遇到了类似的问题。我们通过使用两个工作解决了这个问题。一个人用 spark 计算结果并将其写到“暂存”位置。下一个作业是一个触发的 bash 脚本,它拾取文件并将其连接在一起并传送到必要的位置,然后清理暂存区域。 bash 脚本非常快。不是很优雅,但绝对有效。

标签: apache-spark hdfs apache-spark-sql


【解决方案1】:

(我知道我们可能应该按分区列进行分配,但这是为了提供最简单的示例)

所以看起来我试图简化事情是我出错的地方。如果我DISTRIBUTE BY 实际列而不是人工的1(即DISTRIBUTE BY load_date 或其他),那么它不会创建空文件。为什么?谁知道...

(这也匹配merge-multiple-small-files-in-to-few-larger-files-in-spark 线程上的this 答案)

【讨论】:

    【解决方案2】:

    从 spark 2.4 开始,您可以向查询添加提示以合并和重新分区最终的 Select。例如:

    INSERT OVERWRITE TABLE ... PARTITION (...) 
    SELECT /*+ REPARTITION(5) */ client_id, country FROM mytable;
    

    这将生成 5 个文件。

    在 Spark 2.4 之前,可能会对查询性能产生影响,您可以将 spark.sql.shuffle.partitions 设置为所需文件的数量。

    【讨论】:

    • 感谢 Spark 2.4.3 用户。这就像一个魅力:]
    【解决方案3】:

    这对我来说一直是一个真正令人讨厌的问题,我花了一段时间才解决。

    以下两种方法对我有用:

    1. 作为直线脚本在外部运行:
         set hive.exec.dynamic.partition.mode=nonstrict;
         set hive.merge.mapfiles=true;
         set hive.merge.mapredfiles=true;
         set hive.merge.smallfiles.avgsize=64512000;
         set hive.merge.size.per.task=12992400;
         set hive.exec.max.dynamic.partitions=2048;
         set hive.exec.max.dynamic.partitions.pernode=1024;
    
         <insert overwrite command>
    

    这种方法的问题是这在 pyspark 内部无法以某种方式工作,并且可以作为外部直线脚本从 python 脚本中运行

    1. 使用重新分区

    我发现这个选项非常好。 Repartition(x) 允许将 pyspark 数据帧的记录压缩到“x”个文件中。

    现在,由于表大小各不相同(例如,我不想将具有 1000 万条记录的表重新分区为 1),因此不可能想出一个静态数字“x”来重新分区每个表,我愿意以下。

    -> set an upper threshold for the max number of records a partition should hold 
    I use 100,000 
    
    -> compute x as : 
    x = df.count()/max_num_records_per_partition
    In case the table is partitioned,I use df_partition instead of df...i.e for every set of partition values, i filter df_partition from df; and then compute x from df_partition
    
    -> repartition as:
    df = df.repartition(x)
    In case if the table is partitioned; i use df_partition = df_partition.repartition(x)
    
    -> insert overwrite dataframe 
    

    这种方法对我来说很方便。

    更进一步,表中的列数和用于每列的数据类型可用于创建权重,这些权重可用于更有效地估计给定数据帧的重新分区数。 (例如,与具有 5 列相同类型的数据帧相比,具有 20 列的数据帧将获得更高的权重;与具有类型的 1 列数据帧相比,具有 Map 类型的 1 列数据帧将获得更高的权重布尔值

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-11-15
      • 1970-01-01
      • 2022-07-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多