You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.16 09:05:25