如何在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 DataFrameevaluator.evaluate是作用于整个Spark DataFrame的方法,无法直接处理分组后的子数据集,必须通过pandas UDF封装后才能在分组逻辑中使用
内容的提问来源于stack exchange,提问作者user124123
相关产品推荐
相关产品推荐

