如何在PySpark中统计分组内的连续NA数据?
PySpark统计分组内连续NA的计数
问题场景
现有数据集如下:
| group_num | useage | days | time |
|---|---|---|---|
| 1 | 10 | 20200101 | 1 |
| 1 | 10 | 20200101 | 2 |
| 1 | na | 20200101 | 3 |
| 2 | 30 | 20200102 | 1 |
| 2 | na | 20200102 | 2 |
| 2 | na | 20200102 | 3 |
| 3 | na | 20200105 | 10 |
| 3 | na | 20200105 | 11 |
| 3 | 5 | 20200105 | 12 |
需要按group_num分组,统计useage字段中连续出现的NA次数,遇到非NA值或分组变化时重置计数,最终得到包含na_count字段的结果集。
解决方案
这个需求可以通过PySpark的窗口函数结合条件判断实现,核心思路是先对每个分组内的数据排序,标记出连续NA的分段,最后在分段内计算计数。
步骤1:创建测试DataFrame
先模拟输入数据(实际场景中NA对应PySpark的null):
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window spark = SparkSession.builder.appName("ContinuousNACount").getOrCreate() data = [ (1, 10, "20200101", 1), (1, 10, "20200101", 2), (1, None, "20200101", 3), (2, 30, "20200102", 1), (2, None, "20200102", 2), (2, None, "20200102", 3), (3, None, "20200105", 10), (3, None, "20200105", 11), (3, 5, "20200105", 12) ] df = spark.createDataFrame(data, ["group_num", "useage", "days", "time"])
步骤2:定义窗口并标记连续NA分段
按group_num分组、time排序,通过对比当前行与前一行的NA状态,生成分段标识:
# 定义分组排序窗口 window_group = Window.partitionBy("group_num").orderBy("time") # 标记当前行是否为NA df = df.withColumn("is_na", F.col("useage").isNull()) # 生成分段标识:当前行与前一行NA状态不同时,开启新分段 df = df.withColumn( "segment", F.sum( F.when( F.lag("is_na").over(window_group) != F.col("is_na"), 1 ).otherwise(0) ).over(window_group).cast("int") ) # 处理分组内第一行的分段默认值 df = df.withColumn("segment", F.coalesce("segment", F.lit(0)))
步骤3:计算每个分段内的连续计数
基于分段标识,在group_num+segment的分组内计算行号,再生成最终的na_count:
# 定义分段内排序窗口 window_segment = Window.partitionBy("group_num", "segment").orderBy("time") # 计算分段内的行号 df = df.withColumn("row_num", F.row_number().over(window_segment)) # 生成na_count:非NA行置0,NA行用行号作为连续计数 df = df.withColumn( "na_count", F.when(F.col("is_na"), F.col("row_num")).otherwise(0) ) # 保留目标字段 result_df = df.select("group_num", "useage", "days", "time", "na_count") result_df.show()
运行结果
执行后输出如下:
+---------+------+--------+----+--------+ |group_num|useage| days|time|na_count| +---------+------+--------+----+--------+ | 1| 10|20200101| 1| 0| | 1| 10|20200101| 2| 0| | 1| null|20200101| 3| 1| | 2| 30|20200102| 1| 0| | 2| null|20200102| 2| 1| | 2| null|20200102| 3| 2| | 3| null|20200105| 10| 1| | 3| null|20200105| 11| 2| | 3| 5|20200105| 12| 0| +---------+------+--------+----+--------+
关键说明
- 必须保证每个分组内的数据按
time排序,否则连续NA的判断会出错。 lag函数用于对比前后行的NA状态,是区分连续NA段的核心逻辑。- 分段内的行号直接对应连续NA的计数,非NA行直接置0即可满足需求。
内容的提问来源于stack exchange,提问作者KJH
相关产品推荐
相关产品推荐

