【问题标题】:How to physically partition data to avoid shuffle in Spark SQL joins如何对数据进行物理分区以避免 Spark SQL 连接中的洗牌
【发布时间】:2016-10-24 19:07:22
【问题描述】:

我需要加入 5 个中等大小的表(每个约 80 gb),输入数据约为 800 gb。所有数据都驻留在 HIVE 表中。 我正在使用 Spark SQL 1.6.1 来实现这一点。 加入需要 40 分钟才能完成 --num-executors 20 --driver-memory 40g --executor-memory 65g --executor-cores 6。所有连接都是排序合并外连接。也看到很多洗牌发生。

我将 hive 中的所有表分桶到相同数量的桶中,以便在首先加载数据本身时来自所有表的相似键将转到相同的 spark 分区。但似乎 spark 不理解分桶。

有没有其他方法我可以在 Hive 中对数据进行物理分区和排序(没有部分文件),以便 spark 在从 hive 本身加载数据时知道分区键,并在同一个分区中进行连接,而不需要对数据进行混洗?这将避免在从 hive 加载数据后进行额外的重新分区。

【问题讨论】:

    标签: apache-spark-sql


    【解决方案1】:

    首先,Spark Sql 1.6.1 还不支持 hive 存储桶。 因此,在这种情况下,我们只剩下 Spark 级别的操作,以确保所有表在加载数据时都必须转到相同的 spark 分区。 Spark API 提供了 repartition 和 sortWithinPartitions 来实现相同的目的。例如

    val part1 = df1.repartition(df1("key1")).sortWithinPartitions(df1("key1"))

    以同样的方式,您可以为剩余的表进行几代分区,并将它们连接到在分区内排序的键上。

    这将使连接“无随机播放”操作,但会带来大量计算成本。缓存数据帧(您可以对新创建的分区进行缓存操作)如果该操作将在后续时间执行,则性能会更好。希望这有帮助。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-01-07
      • 1970-01-01
      • 1970-01-01
      • 2019-07-01
      • 2018-06-18
      • 2016-11-28
      相关资源
      最近更新 更多