【发布时间】: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)
当前火花配置:
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