PySpark 3.1.2使用grouping函数与Oracle行为不一致报错咨询
PySpark 3.1版本cube操作grouping函数报错解决方案
这个报错是PySpark 3.1.2及更早版本的已知缺陷,和CUBE/ROLLUP操作的列解析逻辑有关:
- Oracle等传统数据库对
grouping()函数的列引用是按逻辑名称匹配,只要列属于分组维度,就可以正常识别 - PySpark 3.1及更早版本执行
cube()/rollup()时,会给所有分组维度生成带新表达式ID的列,而grouping()函数仍然会匹配原始表的列ID,就会出现你报错中的ID不匹配问题:原始roll列ID为roll#19400,分组后生成的roll列ID为roll#19572,二者无法对应因此抛出异常。
方案1:版本升级
直接升级到PySpark 3.2及以上版本即可,这个缺陷在3.2版本的修复补丁中已经被解决,你当前的SQL不需要做任何修改就可以正常执行。
方案2:SQL改写(不升级版本可用)
如果暂时无法升级Spark版本,可以用两种改写方式规避这个问题:
改写方式1:拆分group by逻辑
把不需要参与cube的列单独放到group by中,避免roll列被重新生成ID:
spark.sql(''' select roll, grouping(roll), grouping(grade), grouping(subject), count(*) from dfview group by roll, cube(grade,subject) having grouping(roll) = 0 and count(*) > 1 ''').show()
改写方式2:用grouping_id()函数替代
通过grouping_id()的位运算实现相同的过滤逻辑,避开单个grouping()函数的列匹配问题:
spark.sql(''' select roll, grouping(roll), grouping(grade), grouping(subject), count(*) from dfview group by cube(roll,grade,subject) having (grouping_id(roll,grade,subject) & 4) = 0 and count(*) > 1 ''').show()
注:
grouping_id按参数从左到右生成二进制位,roll是第一个参数对应最高位4(二进制100),按位与结果为0即代表grouping(roll)=0。
内容的提问来源于stack exchange,提问作者Anirban Chakraborty
相关产品推荐
相关产品推荐

