PySpark DataFrame窗口计数异常:如何实现预期计数结果?
修正PySpark窗口计数问题
你的问题源于原代码使用的count("B").over(W)是累计非空计数,窗口按A分组、C正序时,每行的cnt是从组内第一行到当前行的B非空值数量,因此会出现计数从1开始、且每行是累计到当前行的数值这两个不符合预期的情况。以下是两种针对性的修正方案:
方案一:倒序窗口 + 序号偏移
通过将窗口按C倒序排序,再用row_number()生成序号后减1,直接得到从0开始、正序第一行对应组内最大计数的结果:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 按A分组,C倒序排序的窗口 W = Window.partitionBy("A").orderBy(F.col("C").desc()) main_df = main_df.withColumn("cnt", F.row_number().over(W) - 1)
比如组内按C正序有3行,倒序后序号为1、2、3,减1后得到2、1、0,正好匹配你预期的“正序第一行取组内最后一个计数、计数从0开始”的需求。
方案二:组内总数减去正序序号
先计算每组的总条数,再用总条数减去正序的行号,逻辑更直观:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 计算每组总条数的窗口(无需排序) window_total = Window.partitionBy("A") # 按C正序排序的窗口 window_order = Window.partitionBy("A").orderBy(F.col("C")) main_df = main_df.withColumn("total", F.count("B").over(window_total)) \ .withColumn("cnt", F.col("total") - F.row_number().over(window_order))
如果需要和原代码一致,仅统计B非空的行数,这里的count("B")会自动忽略B为null的行;若要统计所有行(包括B为null),替换为count("*")即可。
内容的提问来源于stack exchange,提问作者Passive_coder
相关产品推荐
相关产品推荐

