【问题标题】:Apache Spark, range-joins, data skew and performanceApache Spark、范围连接、数据倾斜和性能
【发布时间】:2021-08-12 21:20:25
【问题描述】:

我有以下 Apache Spark SQL 连接谓词:

t1.field1 = t2.field1 and t2.start_date <= t1.event_date and t1.event_date < t2.end_date

数据:

t1 DataFrame have over 50 millions rows
t2 DataFrame have over 2 millions rows

t1 DataFrame 中的几乎所有t1.field1 字段都具有相同的值(null)。

目前,由于数据倾斜,Spark 集群在单个任务上挂起超过 10 分钟以执行此连接。此时只有一名工人和该工人的一项任务在工作。所有其他 9 名工人都处于空闲状态。如何改进这种连接,以便将这一特定任务的负载分配到整个 Spark 集群?

【问题讨论】:

    标签: scala apache-spark apache-spark-sql


    【解决方案1】:

    我假设你正在做内部连接。

    可以按照以下步骤来优化连接 - 1.加入前我们可以根据最小或最大的start_date、event_date、end_date过滤掉t1和t2。它会减少行数。

    1. 检查 t2 数据集是否有 field1 的空值,如果没有,则可以根据 notNull 条件过滤连接 t1 数据集。它会减小 t1 的大小

    2. 如果您的作业只获得比可用执行器少的执行器,那么您的分区数就会减少。只需对数据集重新分区,设置一个最佳数字,这样就不会出现大量分区,反之亦然。

    3. 您可以通过查看任务执行时间来检查分区是否正确(无偏斜),应该类似。

    4. 检查较小的数据集是否可以放入执行器内存中,可以使用broadcast_join。

    您可能想阅读 - https://github.com/vaquarkhan/Apache-Kafka-poc-and-notes/wiki/Apache-Spark-Join-guidelines-and-Performance-tuning

    【讨论】:

      【解决方案2】:

      如果 t1 中几乎所有的行都有 t1.field1 = null,并且 event_date 行是数字的(或者您将其转换为时间戳),您可以先使用Apache DataFu 进行范围连接,然后过滤掉 t1.field1 != t2.field1 所在的行。

      范围连接代码如下所示:

      t1.joinWithRange("event_date", t2, "start_date", "end_date", 10)
      

      最后一个参数 - 10 - 是减少因子。正如Raphael Roth 在他的回答中所建议的那样,这确实进行了分桶。

      您可以在the blog post introducing DataFu-Spark 中查看此类远程连接的示例。

      完全披露 - 我是 DataFu 的成员并撰写了博客文章。

      【讨论】:

        【解决方案3】:

        我假设 spark 已经在t1.field1 上推送了非空过滤器,您可以在说明计划中验证这一点。

        我宁愿尝试创建一个可用作等连接条件的附加属性,例如通过桶。例如,您可以创建 month 属性。为此,您需要在t2 中枚举months,这通常使用UDF 完成。有关示例,请参阅此 SO 问题:How to improve broadcast Join speed with between condition in Spark

        【讨论】:

          猜你喜欢
          • 2016-12-21
          • 2015-07-08
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2021-04-23
          • 2017-08-13
          相关资源
          最近更新 更多