【问题标题】:What is an optimized way of joining large tables in Spark SQL在 Spark SQL 中连接大表的优化方法是什么
【发布时间】:2016-06-15 18:01:04
【问题描述】:

我需要使用 Spark SQL 或 Dataframe API 连接表。需要知道实现它的优化方式。

场景是:

  1. Hive 中的所有数据都以 ORC 格式(基本数据框和参考文件)存在。
  2. 我需要将从 Hive 读取的一个 Base 文件(Dataframe)与 11-13 个其他参考文件连接起来,以创建一个大的内存结构(400 列)(大约 1 TB 大小)

实现这一目标的最佳方法是什么?如果有人遇到类似问题,请分享您的经验。

【问题讨论】:

    标签: apache-spark apache-spark-sql


    【解决方案1】:

    我对如何优化连接的默认建议是:

    1. 如果可以,请使用广播连接(请参阅this notebook)。从您的问题来看,您的表似乎很大,广播连接不是一种选择。

    2. 考虑使用一个非常大的集群(您可能认为它更便宜)。现在 250 美元(2016 年 6 月)在 EC2 现货实例市场上购买了大约 24 小时的 800 个内核、6Tb RAM 和许多 SSD。在考虑大数据解决方案的总成本时,我发现人类往往会大大低估自己的时间。

    3. 使用相同的分区程序。请参阅 this question 了解有关联合分组连接的信息。

    4. 如果数据量很大和/或您的集群无法增长以至于即使上述 (3) 导致 OOM,请使用两遍方法。首先,重新分区数据并使用分区表 (dataframe.write.partitionBy()) 进行持久化。然后,在一个循环中连续连接子分区,“附加”到同一个最终结果表。

    旁注:我在上面说“附加”是因为在生产中我从不使用SaveMode.Append。它不是幂等的,这是一件危险的事情。我在分区表树结构的子树中使用SaveMode.Overwrite。在 2.0.0 和 1.6.2 之前,您必须删除 _SUCCESS 否则元数据文件或动态分区发现会阻塞。

    希望这会有所帮助。

    【讨论】:

    • 谢谢你会想到这些。
    • 只是好奇,这台6TB+800核的一台服务器是AWS上的吗?
    • 这是一个集群,而不是单个服务器。是的,在 AWS 上。
    • @Sim。如果你能用一个例子解释你的第四点,我将不胜感激(在循环中连续加入子分区,“附加”到同一个最终结果表)。您能否以两个数据框为例以及如何循环遍历来详细说明一下
    • @vikrantrana 我没有足够的带宽来做这个,但基本的想法很简单:一个大的 USING 连接可以分解成一个连接列空间子集的较小连接的联合。您可以通过对连接的左侧和右侧应用一致的分区来构建子集。例如,如果您要加入一个整数 ID,您可以按 ID 对某个数字取模进行分区,例如 df.withColumn("par_id", id % 256).repartition(256, 'par_id).write.partitionBy("par_id")... 然后迭代 persisted.select('par_id).distinct.collect 加入每个分区 + 再次持久化。然后联合。
    【解决方案2】:

    对源分区使用哈希分区或范围分区,或者如果您更了解连接字段,您可以编写自定义分区。分区将有助于避免在连接期间重新分区,因为跨表的同一分区的火花数据将存在于同一位置。 ORC 肯定会帮助这项事业。 如果这仍然导致溢出,请尝试使用比磁盘更快的 tachyon

    【讨论】:

      【解决方案3】:

      Spark 使用SortMerge joins 加入大表。它包括对两个表上的每一行进行散列,并将具有相同散列的行洗牌到同一个分区中。那里的键在两边都排序,并应用了 sortMerge 算法。据我所知,这是最好的方法。

      要大幅加快您的 sortMerges,请将您的大型数据集编写为具有预存储和预排序选项(相同数量的分区)的 Hive 表,而不是平面 parquet 数据集。

      tableA
        .repartition(2200, $"A", $"B")
        .write
        .bucketBy(2200, "A", "B")
        .sortBy("A", "B")   
        .mode("overwrite")
        .format("parquet")
        .saveAsTable("my_db.table_a")
      
      
      tableb
        .repartition(2200, $"A", $"B")
        .write
        .bucketBy(2200, "A", "B")
        .sortBy("A", "B")    
        .mode("overwrite")
        .format("parquet")
        .saveAsTable("my_db.table_b")
      

      与收益相比,编写预存储/预排序表的开销成本适中。

      默认情况下,底层数据集仍将是 parquet,但 Hive 元存储(可以是 AWS 上的 Glue 元存储)将包含有关表结构的宝贵信息。因为所有可能的“可连接”行都位于同一位置,所以 Spark 不会对预先分桶的表进行洗牌(大大节省了!)并且不会对表分区中预先排序的行进行排序。

      val joined = tableA.join(tableB, Seq("A", "B"))
      

      查看预分桶和不分桶的执行计划。

      这不仅可以在您的联接过程中为您节省大量时间,还可以在没有 OOM 的情况下在相对较小的集群上运行非常大的联接。在亚马逊,我们大部分时间都在生产中使用它(仍有少数情况下不需要它)。

      要了解有关预分拣/预排序的更多信息:

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2021-12-06
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2015-01-16
        • 1970-01-01
        相关资源
        最近更新 更多