【问题标题】:Slow join in pyspark, tried repartition缓慢加入pyspark,尝试重新分区
【发布时间】:2021-08-30 15:35:57
【问题描述】:

我正在尝试在 Spark 3 上左连接 2 个表,其中包含 1700 万行(事件)和 400M 行(详细信息)。 拥有 1 + 15 x 64 核实例的 EMR 集群。 (r6g.16xlarge 尝试使用类似的 r5a) 源文件是从 S3 加载的未分区拼花。

这是我用来加入的代码:

join = (
    broadcast(events).join(
        details,
        [
            details["a"] == events["a2"],
            (unix_timestamp(events["date"]) - unix_timestamp(details["date"])) / 3600
            > 5,
        ],
        "left",
    )
).drop("a")

join.checkpoint()

我用这个来分区:

executors = 15 * 64 * 3  # 15 instances, 64 cores, 3 workers per core

所以我尝试了:

details = details.repartition(executors, "a")

details = details.withColumn("salt", (rand(seed=42) * nSaltBins).cast("int"))
details = details.repartition(executors, "salt")

在这两种情况下,90% 的工作人员会在大约 5-10 分钟内结束,其余的工作人员会持续很长时间(50 分钟以上),绿线长,日志上没有内存或磁盘错误。

分区后有一点偏差(所有分区在 180k 和 160k 行之间),处理器时间超过 50 分钟。

知道我可以监督什么吗?看了一大堆帖子,还是觉得绿线(工人时间)应该更近一些,它们都是同时开始的,不是在等待工人结束。

谢谢!

---编辑--- 已删除广播

在作业 11,第 17 阶段,它在 2 分钟内完成 974/1000,30 分钟后仍然在 993/1000,上一步使用加盐分区(由 executors 变量给出),速度非常快。

执行计划:

Using 17906254 events
== Physical Plan ==
AdaptiveSparkPlan (13)
+- Project (12)
   +- SortMergeJoin LeftOuter (11)
      :- Sort (4)
      :  +- Exchange (3)
      :     +- Project (2)
      :        +- Scan parquet  (1)
      +- Sort (10)
         +- Exchange (9)
            +- Exchange (8)
               +- Project (7)
                  +- Filter (6)
                     +- Scan parquet  (5)

25% 的 2 小时及以上的示例是剩余 1 个执行者

当前火花配置:

spark = SparkSession.builder.appName('Test').config("spark.driver.memory", "108g").config(
        "spark.executor.instances", "59").config("spark.executor.memoryOverhead", "13312").config(
        "spark.executor.memory", "108g").config("spark.executor.cores", "15").config("spark.driver.cores", "15").config(
        "spark.default.parallelism", "1770").config("spark.sql.adaptive.enabled", "true").config(
        "spark.sql.adaptive.skewJoin.enabled", "true").config("spark.sql.shuffle.partitions", "885").getOrCreate()

【问题讨论】:

  • 也尝试重新分区事件,因为它可能是一个排序合并连接。
  • 你应该放弃广播
  • 添加执行计划以发布,删除广播@AdibP 谢谢!
  • @hagarwal 只是为了理解,为什么这会有所帮助?添加执行计划,谢谢!
  • @hagarwal 它似乎在这方面工作得更好一些,但后来也卡住了,我通过盐渍和另一个通过连接字段进行了分区(非常倾斜)谢谢!

标签: apache-spark pyspark apache-spark-sql


【解决方案1】:

您的问题看起来像是一个很好的倾斜联接案例,其中某些分区将获得比其他分区更多的数据,从而减慢整个工作。

在您加入之前重新分区您的数据框将无济于事,因为 SortMergeJoin 操作将在您的加入键上再次重新分区以处理加入

由于您使用的是 Spark 3,因此您应该支持 automatic skewJoin management

要使用它,请确保同时拥有spark.sql.adaptive.enabled=true(在标准 Spark 发行版中默认为 false)和 spark.sql.adaptive.skewJoin.enabled=true

如果您不能使用自动 skewJoin 优化,您可以通过以下方式手动修复它:

  • 将小数据集复制 N 次
n = 10   # Chose an appropriate amount based on skewness
skewedEvents = events.crossJoin(spark.range(0,n).withColumnRenamed("id","eventSalt"))
  • 使用介于 0 和 N 之间的随机列值作为大型数据集的种子
import pyspark.sql.functions as f

skewedDetails = details.withColumn("detailSalt", (f.rand() * n).cast("int"))
  • 在加入键中使用盐加入,然后删除盐
joined = skewedEvents.join(skewedDetails,[         [
            skewedDetails["a"] == skewedEvents["a2"],
            skewedDetails["detailSalt"] == skewedEvents["eventSalt"],
            (unix_timestamp(skewedEvents["date"]) - unix_timestamp(skewedDetails["date"])) / 3600
            > 5,
        ],
        "left")\
        .filter("a is not null or (a is null and eventSalt = 0)")\
        .drop("a").drop("eventSalt").drop("detailSalt")

请注意,您可能还需要验证您的查询连接条件,因为 UI 显示,在详细信息上处理了 3.33 亿行,在事件上处理了 1700 万行,您生成了超过 50 亿的输出行,因此您可以匹配更多您认为的行您的加入条件。

【讨论】:

  • 添加了根据您的输入更新的当前 spark 配置
  • 更新!现在结束了,太好了,CustomShuffleReader 只是在“合并”,所以将分区设置为 spark.sql.shuffle.partitions。 SortMergeJoin 的输出显示 270 亿行,最后 3 个执行程序增加到 44 行(数字正在变化,因为我在连接比较中犯了错误,现在已修复),我应该尝试加盐吗?不想惹AQE。谢谢!!!
  • 对于 AQE,您通常希望将 spark.sql.shuffle.partitionsspark.sql.adaptive.coalescePartitions.initialPartitionNum 设置为您的数据量的高值,并让 AQE 合并太小的分区。您可以查看最大任务的行数,以查看偏差是否仍然足够显着以保证手动加盐。您可能还想将执行程序调整为 5 核/31G 堆内存,您可能会获得一些额外的性能提升
  • 不会重复skewedEvents N 次并按照建议进行 outer 连接,从而为每个事件产生 N-1 额外行吗?使用内部连接不会有问题,但在这种情况下,您需要事后进行某种“清理”。
  • 在大数据帧上,salt 用于将键随机拆分为 N 个拆分(从而将较大的分区大小减少 N 倍)。要使用这些拆分,您需要加入 (key + salt),因此您还需要在小桌子一侧加盐。由于所有拆分都可能连接到小表中的任何行,因此不能随机分配盐,因此我们通过复制每个盐值的所有行来分配它。最终过滤器是必需的,因为您正在对小型数据集进行左外连接,如果在大型数据集中找不到键,您的输出行将重复 N 次。
【解决方案2】:

broadcast() 用于在每个执行器上缓存数据(而不是在每个任务中发送数据),但它在处理大量数据时效果不佳。在这里,1700 万行似乎有点太多了。

如果源数据的分区未针对连接进行优化,则在连接之前对源数据进行预分区也会有所帮助。您需要围绕用于连接的列进行分区。通常应该根据数据的使用方式对数据进行分区。

【讨论】:

  • 嗨!如之前的 cmets 所述,删除广播没有大的改进,如果按连接上使用的 id 进行分区,结果分区会变得非常倾斜:(
猜你喜欢
  • 2016-02-23
  • 1970-01-01
  • 2019-06-29
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-02-04
  • 2019-11-07
  • 1970-01-01
相关资源
最近更新 更多