如何在Spark DataFrame窗口中运行自定义异常值检测函数?
按小时区间检测Spark DataFrame异常值的实现方法
问题说明
现有包含pressure和Timestamp的Spark DataFrame,已新增date_hour列用于按小时分组,需要基于每个小时区间的数据检测异常值,而非全数据集。目前已有针对全数据集的异常值检测函数,但无法直接应用到窗口分组中。
现有全数据集异常值检测函数
def outlierDetection(df): inter_quantile_range = df.approxQuantile("pressure",[0.20,0.80],relativeError=0) Q1=inter_quantile_range[0] Q3=inter_quantile_range[1] inter_quantile_diff = Q3 - Q1 minimum_Q1 = Q1 - 1.5 * inter_quantile_diff maximum_Q3 = Q3 + 1.5 * inter_quantile_diff df= df.withColumn("isOutlier",F.when((df["pressure"] > maximum_Q3) | (df["pressure"] < minimum_Q1), 1).otherwise(0)) return df
按小时分组的异常值检测实现
由于approxQuantile无法直接在窗口函数中使用,推荐通过分组聚合计算小时级分位数,再关联原数据的方式实现:
步骤1:计算每个小时区间的分位数及异常值阈值
from pyspark.sql import functions as F # 分组计算每个date_hour的Q1(20%)和Q3(80%) hourly_quantiles = df.groupBy("date_hour").agg( F.expr("percentile_approx(pressure, 0.20, 10000)").alias("Q1"), F.expr("percentile_approx(pressure, 0.80, 10000)").alias("Q3") ) # 计算每个小时的异常值上下限 hourly_thresholds = hourly_quantiles.withColumn( "IQR", F.col("Q3") - F.col("Q1") ).withColumn( "min_threshold", F.col("Q1") - 1.5 * F.col("IQR") ).withColumn( "max_threshold", F.col("Q3") + 1.5 * F.col("IQR") )
注:
percentile_approx是Spark的近似分位数函数,第三个参数是精度,值越大结果越准确,这里设为10000足够满足大部分场景需求。
步骤2:关联原数据并标记异常值
# 将阈值表与原DataFrame关联 result_df = df.join(hourly_thresholds, on="date_hour", how="left") # 标记异常值 result_df = result_df.withColumn( "isOutlier", F.when( (F.col("pressure") > F.col("max_threshold")) | (F.col("pressure") < F.col("min_threshold")), 1 ).otherwise(0) ).drop("Q1", "Q3", "IQR", "min_threshold", "max_threshold") # 可选:删除中间计算列
替代方案:使用窗口函数实现(适用于Spark 3.0+)
如果使用Spark 3.0及以上版本,可以结合percent_rank窗口函数计算分位数,但这种方式是基于排序后的相对位置,结果与approxQuantile略有差异:
from pyspark.sql import Window w = Window.partitionBy("date_hour").orderBy("pressure") # 计算每个值在小时分组内的百分位排名 df_with_rank = df.withColumn("p_rank", F.percent_rank().over(w)) # 提取Q1和Q3的阈值(这里取百分位20%和80%对应的压力值) hourly_ranks = df_with_rank.groupBy("date_hour").agg( F.max(F.when(F.col("p_rank") <= 0.20, F.col("pressure"))).alias("Q1"), F.min(F.when(F.col("p_rank") >= 0.80, F.col("pressure"))).alias("Q3") ) # 后续步骤同方案一:计算阈值、关联原数据、标记异常值
关键说明
- 避免直接在窗口中使用
approxQuantile:该函数是Action操作,无法作为窗口函数的一部分执行,会导致性能问题或错误。 - 分组聚合+关联的方式是最高效且易维护的实现方案,适合大规模数据集。
内容的提问来源于stack exchange,提问作者FlashBang
相关产品推荐
相关产品推荐

