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

如何在PySpark DataFrame分组中应用pyspark.ml的AUC函数

按分组计算PySpark中AUC/PR曲线下面积的解决方案

问题描述

需要将pyspark.ml.evaluation.BinaryClassificationEvaluator的AUC/AUPR计算逻辑应用到PySpark DataFrame的分组数据中(按marketplace_id分组)。已能对整个DataFrame计算指标,但分组计算时报错:

ValueError: Invalid udf: the udf argument must be a pandas_udf of type GROUPED_MAP

原数据读取及评估器初始化代码:

predictionAndSprScoreLabels_sda = spark.table('sda_model').filter(
    (col('region_id') == region_id) & 
    (col('weekend_day') == to_date(lit(execParams['date']), 'yyyy-MM-dd'))
).select( 
    col('attributed_units_sold').cast('double').alias('attributed_units_sold'),
    col('sdar_score').cast('double').alias('sdar_score'),
    col('marketplace_id').cast('double').alias('marketplace_id')
)

evaluator = BinaryClassificationEvaluator(
    labelCol="attributed_units_sold", 
    rawPredictionCol="sdar_score"
)

原分组报错代码:

predictionAndSprScoreLabels_sda=  predictionAndSprScoreLabels_sda.withColumn(
    "auc", 
    predictionAndSprScoreLabels_sda.groupBy("marketplace_id").apply(
        evaluator.evaluate(predictionAndSprScoreLabels_sda, {evaluator.metricName: "areaUnderROC"})
    )
)
predictionAndSprScoreLabels_sda=  predictionAndSprScoreLabels_sda.withColumn(
    "aupr", 
    predictionAndSprScoreLabels_sda.groupBy("marketplace_id").apply(
        evaluator.evaluate(predictionAndSprScoreLabels_sda, {evaluator.metricName: "areaUnderPR"})
    )
)

解决方案

使用GROUPED_MAP类型的pandas UDF封装评估逻辑,对每个分组单独计算AUC和AUPR,再将结果合并回原DataFrame。

完整代码示例

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, lit, to_date, pandas_udf
from pyspark.sql.types import StructType, StructField, DoubleType
from pyspark.ml.evaluation import BinaryClassificationEvaluator

# 1. 读取数据(保留原逻辑)
predictionAndSprScoreLabels_sda = spark.table('sda_model').filter(
    (col('region_id') == region_id) & 
    (col('weekend_day') == to_date(lit(execParams['date']), 'yyyy-MM-dd'))
).select( 
    col('attributed_units_sold').cast('double').alias('attributed_units_sold'),
    col('sdar_score').cast('double').alias('sdar_score'),
    col('marketplace_id').cast('double').alias('marketplace_id')
)

# 2. 定义分组计算的结果Schema
result_schema = StructType([
    StructField("marketplace_id", DoubleType(), nullable=True),
    StructField("auc", DoubleType(), nullable=True),
    StructField("aupr", DoubleType(), nullable=True)
])

# 3. 定义GROUPED_MAP类型的pandas UDF
@pandas_udf(result_schema, functionType=pandas_udf.GROUPED_MAP)
def calculate_grouped_metrics(df):
    # 获取当前活跃的SparkSession,将分组的pandas DataFrame转为Spark DataFrame
    spark = SparkSession.getActiveSession()
    group_spark_df = spark.createDataFrame(df)
    
    # 初始化评估器
    evaluator = BinaryClassificationEvaluator(
        labelCol="attributed_units_sold", 
        rawPredictionCol="sdar_score"
    )
    
    # 计算当前分组的AUC和AUPR
    auc = evaluator.evaluate(group_spark_df, {evaluator.metricName: "areaUnderROC"})
    aupr = evaluator.evaluate(group_spark_df, {evaluator.metricName: "areaUnderPR"})
    
    # 返回分组ID及对应的指标(去重保证每个分组只返回一行)
    return df[["marketplace_id"]].drop_duplicates().assign(auc=auc, aupr=aupr)

# 4. 按marketplace_id分组计算指标
grouped_metrics = predictionAndSprScoreLabels_sda.groupBy("marketplace_id").apply(calculate_grouped_metrics)

# 5. 将指标合并回原DataFrame(可选,根据需求决定是否保留原数据行)
final_result_df = predictionAndSprScoreLabels_sda.join(grouped_metrics, on="marketplace_id", how="left")

错误原因说明

原代码直接在groupBy.apply中调用evaluator.evaluate不符合API要求:

  • groupBy.apply仅支持GROUPED_MAP类型的pandas UDF作为参数,该UDF需要接收分组后的pandas DataFrame并返回结构匹配的pandas DataFrame
  • evaluator.evaluate是作用于整个Spark DataFrame的方法,无法直接处理分组后的子数据集,必须通过pandas UDF封装后才能在分组逻辑中使用

内容的提问来源于stack exchange,提问作者user124123

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 23:01:14