PySpark 2.2迁移至3.1.2后collect_set返回全Null问题咨询
PySpark 2.2 迁移3.1.2:collect_set窗口函数返回全Null问题解析
问题根源:窗口函数默认行为变更
PySpark 2.2与3.1.2中,聚合类窗口函数的默认窗口帧规则存在差异:
- 在2.2版本中,使用
collect_set作为窗口函数时,默认覆盖整个分区(等价于显式指定rows between unbounded preceding and unbounded following),会收集分区内所有非Null的唯一值。 - 在3.1.2版本中,若窗口定义未指定
ORDER BY,聚合类窗口函数的默认窗口帧变为range between unbounded preceding and current row——这意味着collect_set仅会收集从分区起始行到当前行的唯一值。如果分区内行的顺序导致当前行及之前无有效非Null值,就会返回包含Null的集合,表现为全Null结果。
解决方案
显式指定窗口帧覆盖整个分区,修改后的SQL如下:
sqlContext.sql( select collect_set({column}) over(partition by {partition_column} rows between unbounded preceding and unbounded following) as {col} )
额外检查点
- 确认
{column}字段本身存在非Null的地区数据,排除数据源本身的问题。 - PySpark 3.x中
collect_set仍会自动排除集合中的Null值(与2.x一致),仅当收集范围内无有效非Null值时才会返回含Null的集合。
内容的提问来源于stack exchange,提问作者baka yaro
相关产品推荐
相关产品推荐

