【问题标题】:how to manually run pyspark's partitioning function for debugging如何手动运行pyspark的分区功能进行调试
【发布时间】:2021-04-17 07:06:11
【问题描述】:

我遇到了性能问题,在检查了 spark 的 Web UI 后,我发现存在严重的数据偏斜问题:

我已经尝试通过列的多个组合“分发”和“重新分区”但没有运气,所以我试图调试 spark 是如何对数据集进行分区的(为了修复它),有什么方法可以手动运行用于创建列的分区函数?基本上我正在尝试做类似的事情:

df = df.withColumn("assigned_partition", partitioning_function())
df_grouped  = df.groupby("assigned_partition").count()

这样我就可以确定偏斜的模式或原因。

注意:这是在查询 hive 表之后,所以我知道偏度不是由于任何 spark 逻辑或计算造成的。

【问题讨论】:

    标签: apache-spark pyspark


    【解决方案1】:

    我发现可以使用:spark_partition_id() ,然后我创建了一个直方图来分析 count() 分布,发现它看起来正常且没有异常值,因此偏度与数据集的分区方式无关。

    test_df = df.select(spark_partition_id().alias("partitionId"))
    
    test_df.groupBy("partitionId").count().orderBy(col("count").desc()).select("count").toPandas().plot.hist()
    
    plt.show()
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2015-03-29
      • 1970-01-01
      • 2019-05-13
      • 2012-12-29
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-04-12
      相关资源
      最近更新 更多