PySpark DataFrame中计算每行score均值与review众数报错求助
问题描述
现有如下结构的PySpark DataFrame:
+----+--------------------------+--------------+ |id | score |review | +----+--------------------------+--------------+ |1 |[83.52, 81.79, 84, 75] |[P,N,P,P] | |2 |[86.13, 85.48] |[N,N,N,P] | +----+--------------------------+--------------+
Schema为:
root |--id: int (nullable = false) |--score: array (nullable = true) | |-- element: string (containsNull = true) |--review : array (nullable = true) | |-- element: string (containsNull = true)
需要为每行计算score列数组的均值(生成float类型的scoreMean列)和review列数组的众数(生成string类型的reviewMode列),期望结果如下:
+----+--------------------------+--------------+---------+----------+ |id | score |review |scoreMean|reviewMode| +----+--------------------------+--------------+---------+----------+ |1 |[83.52, 81.79, 84, 75] |[P,N,P,P] |81.08 |P | |2 |[86.13, 85.48] |[N,N,N,P] |85.81 |N | +----+--------------------------+--------------+---------+----------+
尝试使用UDF实现,但执行show或collect时出现计算失败错误,代码如下:
import statistics import pyspark.sql.functions as F #define UDFs def mean_udf(data): if len(data) == 0: return None data_float = [eval(i) for i in data] return statistics.mean(data_float) def mode_udf(data): if len(data) == 0: return None return statistics.mode(data) #register the UDFs mean_func = F.udf(mean_udf) mode_func = F.udf(mode_udf) #apply UDFs df= (df.withColumn("scoreMean", mean_func(F.col("score"))) .withColumn("reviewMode", mode_func(F.col("review"))) )
错误原因分析
eval函数的问题:分布式环境中eval可能因序列化报错,且直接执行字符串有安全风险,若数组元素为空或格式错误,eval会抛出异常。statistics.mode的局限性:当数组存在多个众数时,该函数会抛出StatisticsError导致任务失败。- UDF未指定返回类型:默认UDF返回
StringType,但scoreMean需要float类型,类型不匹配可能引发问题。 - 空值处理不全:若
score或review列为null,访问len(data)会抛出TypeError。
解决方案
方案一:修正UDF实现
针对上述问题,优化后的UDF如下:
import statistics from collections import Counter import pyspark.sql.functions as F from pyspark.sql.types import FloatType, StringType def safe_mean_udf(data): if not data: # 处理null或空数组 return None try: data_float = [float(i) for i in data if i is not None] # 跳过数组中的null元素 if not data_float: return None return round(statistics.mean(data_float), 2) # 保留两位小数匹配期望结果 except (ValueError, TypeError): return None def safe_mode_udf(data): if not data: # 处理null或空数组 return None non_null_data = [item for item in data if item is not None] if not non_null_data: return None count = Counter(non_null_data) max_count = max(count.values()) # 取第一个出现的众数(若有多个) mode = next(k for k, v in count.items() if v == max_count) return mode # 显式指定UDF返回类型 mean_func = F.udf(safe_mean_udf, FloatType()) mode_func = F.udf(safe_mode_udf, StringType()) # 应用UDF df = (df.withColumn("scoreMean", mean_func(F.col("score"))) .withColumn("reviewMode", mode_func(F.col("review"))) )
方案二:使用PySpark内置函数(推荐)
UDF在分布式计算中性能较差,优先使用内置函数实现:
import pyspark.sql.functions as F from pyspark.sql import Window from pyspark.sql.types import FloatType # 计算score均值:先将数组元素转为float,再求平均 df_score_mean = df.withColumn("score_element", F.explode(F.col("score"))) \ .withColumn("score_element", F.col("score_element").cast(FloatType())) \ .groupBy("id", "score", "review") \ .agg(F.round(F.avg("score_element"), 2).alias("scoreMean")) # 计算review众数:统计每个元素出现次数,取次数最多的 window = Window.partitionBy("id", "review").orderBy(F.desc("count")) df_mode = df.withColumn("review_element", F.explode(F.col("review"))) \ .filter(F.col("review_element").isNotNull()) \ .groupBy("id", "review", "review_element") \ .agg(F.count("review_element").alias("count")) \ .withColumn("rank", F.row_number().over(window)) \ .filter(F.col("rank") == 1) \ .drop("count", "rank") \ .withColumnRenamed("review_element", "reviewMode") # 合并结果 final_df = df_score_mean.join(df_mode, on=["id", "score", "review"], how="left")
说明:
- 内置函数方案避免了Python UDF的序列化和性能问题,更适合大规模数据处理。
- 众数计算使用窗口函数,确保多个众数时取第一个出现的元素,符合示例需求。
- 均值计算通过
explode展开数组,聚合后求平均并保留两位小数。
内容的提问来源于stack exchange,提问作者karan
相关产品推荐
相关产品推荐

