PySpark改写SQL去重计数查询返回0 与SQL执行结果不一致
PySpark改写SQL计数为0的逻辑排查
问题背景
现有可正常执行的SQL查询,用于统计user表中满足多条件的去重studen_num数量;改写为PySpark代码后执行计数结果恒为0,与SQL返回的正确结果不符,需定位逻辑错误。
原正确SQL代码
select count(distinct(studen_num)) from user where email is null and matnumber is null and (academic_history is null or exists (academic_history, x -> x.grade = "A" and x.level is null)) and graduation_status is null and exists (investigated_result, x -> x.panelseated[0].grade = "1" and x.level is null)
待排查的异常PySpark代码
df_count_student = df.filter(df.email.isNull() & df.matnumber.isNull() & (df.academic_history.isNull() | (array_contains(df.academic_history.grade, "A") & df.academic_history.level.isNull())) & (array_contains(df.investigated_result.panelseated[0].grade, "1") & df.investigated_result.level.isNull())) display(df_count_student(countDistinct())
核心错误点
- 数组exists逻辑实现完全错误
原SQL的exists(数组, x -> 条件)语义是:遍历数组内每一个元素,判断是否存在至少一个元素同时满足所有自定义条件。
现有实现有两个本质问题:一是array_contains(df.academic_history.grade, "A")仅判断数组中是否有元素grade为A,没有和「同元素level为空」做绑定,会出现跨元素误判;二是df.academic_history.level.isNull()是判断整个academic_history数组字段本身是否为空,不是判断数组内元素的level字段为空,和SQL语义完全不符。 - 漏写过滤条件
原SQL包含graduation_status is null的过滤规则,现有PySpark过滤逻辑中完全缺失该条件。 - 聚合统计语法错误
最后统计去重数的写法df_count_student(countDistinct())是无效语法:既没有指定countDistinct作用的字段,也没有正确调用DataFrame的agg聚合方法,且代码末尾缺失右括号,本身就无法正常执行。
正确PySpark实现
直接通过expr复用和原SQL完全一致的数组高阶判断逻辑,避免手动拆解数组字段带来的语义偏差,代码如下:
from pyspark.sql.functions import countDistinct, expr # 过滤条件1:1对齐原SQL逻辑 df_filtered = df.filter( df.email.isNull() & df.matnumber.isNull() & ( df.academic_history.isNull() | expr("exists(academic_history, x -> x.grade = 'A' and x.level is null)") ) & df.graduation_status.isNull() & expr("exists(investigated_result, x -> x.panelseated[0].grade = '1' and x.level is null)") ) # 统计去重学生编号数量 result = df_filtered.agg(countDistinct("studen_num").alias("distinct_student_count")) # 展示结果 result.show()
内容的提问来源于stack exchange,提问作者umagba alex
相关产品推荐
相关产品推荐

