You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.25 13:45:28