PySpark指定列循环计算分位数标记异常值报错如何解决?
PySpark 指定列异常值标记实现方案
错误原因
你遇到的报错核心是列选择语法错误:PySpark DataFrame 选取多列需使用方括号[]而非圆括号(),你写的df(df['cl_Left_yCm','cl_Right_yCm'])相当于把DataFrame对象当做函数调用,因此触发'DataFrame' object is not callable报错。
实现步骤
1. 修正分位数计算逻辑
首先正确选取目标列,计算每列的25分位(Q1)、75分位(Q3):
# 指定需要计算异常值的列 target_cols = ["cl_Left_yCm", "cl_Right_yCm"] # 计算每列的Q1、Q3 bounds = { c: dict( zip(["q1", "q3"], df.approxQuantile(c, [0.25, 0.75], 0)) ) for c in target_cols }
其中approxQuantile第三个参数为相对误差,设为0代表计算精确分位数,数据量较大时可适当调大(如0.01)提升计算效率。
2. 计算异常判定阈值
行业通用的箱线图异常判定规则为:小于Q1 - 1.5*IQR 或大于 Q3 + 1.5*IQR 的值为异常值,其中IQR = Q3 - Q1为四分位距。我们可以提前把阈值存入bounds字典方便后续调用:
for c in bounds: iqr = bounds[c]['q3'] - bounds[c]['q1'] bounds[c]['lower'] = bounds[c]['q1'] - 1.5 * iqr bounds[c]['upper'] = bounds[c]['q3'] + 1.5 * iqr
3. 生成异常标记列
支持两种标记模式,按需选择即可(标记逻辑为1=异常,0=正常):
from pyspark.sql import functions as F # 模式1:给每个目标列生成单独的异常标记列,列名后缀为_isOutlier for c in target_cols: df = df.withColumn( f"{c}_isOutlier", F.when( (F.col(c) < bounds[c]['lower']) | (F.col(c) > bounds[c]['upper']), 1 ).otherwise(0) ) # 模式2:生成总异常标记列,只要任意目标列异常就标记为1 df = df.withColumn( "total_isOutlier", F.greatest(*[F.col(f"{c}_isOutlier") for c in target_cols]) )
完整可运行测试示例
你可以用以下代码在本地验证逻辑正确性:
from pyspark.sql import SparkSession from pyspark.sql import functions as F # 初始化Spark spark = SparkSession.builder.appName("outlier_detect").getOrCreate() # 构造示例数据 data = [ ("a1",132,0,-153),("a1",132,0,-153),("a1",129,0,-153),("a1",129,0,-153), ("a2",129,0,-152),("a2",129,0,-152),("a2",130,0,-152),("a2",130,0,-152),("a2",130,0,-152), ("a3",134,0,-147) ] df = spark.createDataFrame(data, schema=["ID","cl_Left_yCm","cl_Right_xCm","cl_Right_yCm"]) # ------------ 异常检测逻辑 ------------ target_cols = ["cl_Left_yCm", "cl_Right_yCm"] # 计算分位数 bounds = { c: dict(zip(["q1", "q3"], df.approxQuantile(c, [0.25, 0.75], 0))) for c in target_cols } # 计算阈值 for c in bounds: iqr = bounds[c]['q3'] - bounds[c]['q1'] bounds[c]['lower'] = bounds[c]['q1'] - 1.5 * iqr bounds[c]['upper'] = bounds[c]['q3'] + 1.5 * iqr # 打异常标记 for c in target_cols: df = df.withColumn( f"{c}_isOutlier", F.when((F.col(c) < bounds[c]['lower']) | (F.col(c) > bounds[c]['upper']),1).otherwise(0) ) df = df.withColumn("total_isOutlier", F.greatest(*[F.col(f"{c}_isOutlier") for c in target_cols])) # 输出结果 df.show()
内容的提问来源于stack exchange,提问作者jenzen
相关产品推荐
相关产品推荐

