Spark分组DataFrame执行max聚合操作返回错误值的原因排查
问题根因
这个问题的核心原因是 PluginDuration 列的数据类型为字符串类型(StringType),而非数值类型。
- 执行
max聚合时,Spark会按照字符串字典序比较大小:字符'9'的ASCII码远大于'2',因此"9.2875"会被判定为比"20.50027778"更大的值,最终返回错误的最大值结果。 - 你看到的平均值结果符合预期,是因为
avg函数作用于字符串列时,Spark会隐式将可转换为数值的字符串转为Double类型计算,你给出的4条记录的数值平均值恰好和结果一致,刚好掩盖了类型错误的问题。
验证方法
执行以下代码查看DataFrame的字段类型,即可确认问题:
df.printSchema()
如果输出中PluginDuration对应的类型是StringType,即可匹配上述根因。
修复方案
在聚合前先将PluginDuration列强制转换为数值类型(比如DoubleType),再执行分组聚合即可得到正确结果:
from pyspark.sql import functions as F # 先转换字段类型 df = df.withColumn('PluginDuration', F.col('PluginDuration').cast('double')) # 再执行聚合 fdf = df.groupBy('CPID').agg( F.max('PluginDuration').alias('max_duration'), F.avg('PluginDuration').alias('avg_duration') ) fdf.filter('CPID = "AN04773"').show()
内容的提问来源于stack exchange,提问作者Scarface
相关产品推荐
相关产品推荐

