在Databricks用PySpark遇median函数缺失问题,求替代方案
解决Spark 3.2.1中无法使用median函数计算中位数的问题
核心结论
percentile_approx("follower_count", 0.5)完全可以作为中位数的替代方案:
- 中位数本质就是第50百分位数,
percentile_approx传入参数0.5,就是计算第50百分位的近似值 - 大数据集下,近似结果和精确中位数差异极小,且性能远高于精确计算;小数据集可通过调整精度参数让结果更接近精确值
可行解决方案
方案1:使用percentile_approx(推荐,适配大数据量)
这是Spark低版本中替代median的官方标准方式,直接用你尝试的代码即可:
import pyspark.sql.functions as f df = spark.table('global_temp.user').join(spark.table('global_temp.pin'), 'ind') df = df.withColumn('age_groups', f.when(f.col('age').between(18, 24), '18-24') .when(f.col('age').between(25, 35), '25-35') .when(f.col('age').between(35, 50), '35-50') .when(f.col('age') > 50, '50+')) # 计算近似中位数 result_df = df.groupBy('age_groups').agg( f.percentile_approx("follower_count", 0.5).alias("median_follower_count") )
如果需要更高精度的近似结果,可添加第三个精度参数(值越小越精确,默认10000):
f.percentile_approx("follower_count", 0.5, 100000).alias("median_follower_count")
方案2:计算精确中位数(仅适合小数据量)
若数据集规模较小,想要完全精确的中位数,可通过窗口函数排序后取中间值:
import pyspark.sql.functions as f from pyspark.sql.window import Window df = spark.table('global_temp.user').join(spark.table('global_temp.pin'), 'ind') df = df.withColumn('age_groups', f.when(f.col('age').between(18, 24), '18-24') .when(f.col('age').between(25, 35), '25-35') .when(f.col('age').between(35, 50), '35-50') .when(f.col('age') > 50, '50+')) # 按年龄组分段,对关注者数量排序并标记位置 window = Window.partitionBy('age_groups').orderBy('follower_count') df_with_rank = df.withColumn('row_num', f.row_number().over(window))\ .withColumn('total_count', f.count('*').over(Window.partitionBy('age_groups'))) # 根据数据行数奇偶性取中位数(偶数行取中间两数的平均值) result_df = df_with_rank.filter( (f.col('row_num') == f.floor((f.col('total_count') + 1)/2)) | (f.col('row_num') == f.floor((f.col('total_count') + 2)/2)) ).groupBy('age_groups').agg(f.avg('follower_count').alias('median_follower_count'))
注意:该方法需要全量排序,大数据量下性能极差,仅适用于小数据集场景。
内容的提问来源于stack exchange,提问作者j stevenage
相关产品推荐
相关产品推荐

