Scala与Spark SQL对同数据集sum聚合结果数值差异问题咨询
问题现象
对同一视图view_error_log分别使用Scala和Spark SQL执行sum聚合操作得到的结果存在差异:Spark SQL查询返回的sum值远高于Scala查询的输出结果
复现代码
Scala 代码
import org.apache.spark.sql.functions._ val df = spark.sql( "select * from view_error_log") val grpdf = df.groupby("id","error_description","table", "log_date").agg(sum("totalRows").alias("totalRows"))
Spark SQL 代码
Select id,error_description,table,log_date,sum(totalRows) as TotalRows from view_error_log group by id,error_description,table,log_date
可能的原因
- 保留字未转义导致字段读取错误:
table是Spark SQL的保留关键字,提交的SQL语句中直接使用table作为查询字段,没有加反引号转义,会导致Spark SQL误解析该字段,实际分组时使用了非预期的table值,大量原本应该拆分的分组被合并,sum结果偏大。 - 大小写配置不一致:Spark默认
spark.sql.caseSensitive参数为false,如果视图中的分组字段或聚合字段存在大小写拼写差异(例如视图中实际字段为TotalRows,Scala代码中写的是totalRows),Scala代码可能读取到空值,sum时空值被忽略导致结果偏小,而SQL语句中的大小写会被自动匹配到正确字段,计算正常。 - 缓存数据不一致:如果之前操作中对
view_error_log对应的DataFrame做过缓存,Scala代码中读取的是未更新的旧缓存数据,而SQL查询会直接扫描底层最新的表数据,两次计算的数据源本身存在差异。 - 隐式类型转换导致分组逻辑差异:如果分组字段存在类型不统一的情况(例如
log_date字段同时存在字符串类型和日期类型的值),Scala代码中读取时会做严格的类型匹配,不同类型的值会被分到不同组,聚合结果更分散总和更小,而Spark SQL会自动做类型兼容转换,将原本不同类型的相同时间值分到同一组,sum结果更大。 - Scala代码未执行全量计算:Spark是懒执行引擎,如果仅定义了
grpdf没有执行触发全量计算的action算子(例如直接调用show()只返回前20条数据的聚合结果,没有统计全量),得到的只是部分数据的sum值,远小于SQL全量计算的结果。 - 视图过滤条件的时态差异:如果
view_error_log的定义中包含基于当前时间的动态过滤条件(例如where log_date >= current_date()-1),两次查询执行的时间间隔刚好跨了时间分区,SQL查询扫描到的数据量比Scala代码扫描的更多,导致sum结果偏大。
内容的提问来源于stack exchange,提问作者abhishek gaikwad
相关产品推荐
相关产品推荐

