【问题标题】:Efficient joins in spark dataframeSpark数据框中的高效连接
【发布时间】:2018-10-16 17:23:35
【问题描述】:

我正在尝试在同一个 DataFrame 上按顺序加入一个 DataFrame (dfA)。 假设 dfA 有 id_x 和 id_y 列,而 dfB 有 id 列和其他一些列。

我想执行以下操作:

dfA.join(dfB, dfA("id_x") === dfB("id")).join(dfB, dfA("id_y") === dfB("id"))

我可以做任何重新分区或预处理来加快速度吗?

【问题讨论】:

    标签: scala apache-spark join apache-spark-sql distributed-computing


    【解决方案1】:

    您使用的是什么版本的火花? Tuning Spark 是一门艺术,本身就是一个广阔的话题。只是盲目地增加分区数量并不总是有帮助。我建议看看以下地方的线索:

    1. 仔细查看 Spark UI 并分析您的 DAG。瓶颈在哪里?是在等待 CPU、内存、磁盘 IO 吗?洗牌太多?
    2. 您的数据是否存在偏差?很少有任务长时间运行,而大多数任务都很快完成?
    3. 您使用了哪种转换?如果可能,请粘贴您的代码摘录。
    4. Bucketing 是 Spark 中的一项新功能,人们普遍认为它有助于连接。但调查您的 DAG 始终是最好的线索来源。
    5. 同样根据您的代码,您希望在什么情况下使用 dfA("id_x") 和 dfA("id_y") 与 dfB("id") 连接?您可能可以在下面尝试一些东西,而不是在连接条件中使用 OR

      val joinCondition = when($"dfA.id_y".isNull, $"dfA.id_y"===$"dfB.id") .otherwise($"dfA.id_x"===$"dfB.id")

      val dfJoined = dfA.join(dfB, joinCondition)

    请告诉我你的发现。

    【讨论】:

      【解决方案2】:

      您可以在 1 次加入中做到这一点:

      dfA.join(dfB, dfA("id_x") === dfB("id") or dfA("id_y") === dfB("id"))
      

      您也可以使用spark.sql.shuffle.partitions 或尝试广播一个数据帧。在连接之前重新分区不会有帮助,但使用分桶表可能会有所帮助,因为这可以避免在连接期间进行重新分区,例如https://issues.apache.org/jira/browse/SPARK-12394

      【讨论】:

      • 为了其他人的利益,请详细说明“在加入之前重新分区将无济于事。” @Raphael Roth
      • @thebluephantom 我不确定你的意思。 AFAIK 唯一可以提高性能的方法是使用自 spark 2.0 以来的 hives bucketing 功能(请参阅issues.apache.org/jira/browse/SPARK-12394)
      • 否则你会做出非常明确的声明。丹尼尔问题想知道你为什么这么说。
      • @RaphaelRoth 谢谢!我玩过 spark.sql.shuffle.partitions,但我的数量已经很高(4000-8000),而且两个数据帧都是约 3 亿条记录,所以广播并不是真的可行。我将研究蜂巢分桶。
      猜你喜欢
      • 1970-01-01
      • 2017-08-16
      • 1970-01-01
      • 2016-05-14
      • 2021-09-02
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多