小编evy*_*miz的帖子

Spark 窗口函数零偏斜

最近,我在运行 PySpark 作业之一时遇到了问题。在分析 Spark UI 中的阶段时,我注意到运行时间最长的阶段需要 1.2 小时才能运行完,而整个流程运行所需的总时间为 2.5 小时。

SparkUI 阶段选项卡按最长持续时间排序

查看阶段详细信息后,我清楚地发现我面临着严重的数据偏差,导致单个任务运行了整个1.2 小时,而所有其他任务在23 秒内完成。

任务分配显示出非常明显的偏差

总结显示了最长的任务与绝大多数人之间的巨大差异

DAG 显示此阶段涉及窗口函数,它帮助我快速将有问题的区域缩小到几个查询,并找到根本原因 -> 中account使用的列,Window.partitionBy("account")25% 的空值。我没有兴趣计算空帐户的总和,尽管我确实需要涉及的行进行进一步计算,因此我无法在窗口函数之前将它们过滤掉。

这是我的窗口函数查询:

problematic_account_window = Window.partitionBy("account")

sales_with_account_total_df = sales_df.withColumn("sum_sales_per_account", sum(col("price")).over(problematic_account_window))
Run Code Online (Sandbox Code Playgroud)

所以我们找到了罪魁祸首——我们现在能做什么?我们如何解决倾斜和性能问题?

skew apache-spark apache-spark-sql pyspark spark-window-function

1
推荐指数
1
解决办法
828
查看次数