PySpark按Id分组统计连续出现的Minute=1的次数
PySpark实现连续Minute=1序列的次数统计
问题描述
给定如下DataFrame:
| Id | Minute | identifier |
|---|---|---|
| aa | 1 | 1 |
| aa | 1 | 2 |
| aa | 1 | 3 |
| aa | 2(忽略) | 4 |
| aa | 1 | 5 |
| aa | 1 | 6 |
| bb | 1 | 7 |
| bb | 1 | 8 |
| bb | 5(忽略) | 9 |
| bb | 1 | 10 |
需求:按Id分组,统计连续出现的Minute值为1的序列长度,输出每个连续序列对应的Id和次数,期望结果如下:
| Id | Minute |
|---|---|
| aa | 3 |
| aa | 2 |
| bb | 2 |
| bb | 1 |
实现步骤与代码
1. 初始化SparkSession与加载测试数据
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window spark = SparkSession.builder.appName("ConsecutiveMinuteCount").getOrCreate() # 构造测试数据 data = [ ("aa", "1", 1), ("aa", "1", 2), ("aa", "1", 3), ("aa", "2(忽略)", 4), ("aa", "1", 5), ("aa", "1", 6), ("bb", "1", 7), ("bb", "1", 8), ("bb", "5(忽略)", 9), ("bb", "1", 10) ] df = spark.createDataFrame(data, schema=["Id", "Minute", "identifier"])
2. 清洗Minute列并标记目标行
提取Minute列中的纯数字并转为整数,同时标记当前行是否为Minute=1的有效行:
# 提取Minute中的数字部分并转换为整数类型 df_clean = df.withColumn("Minute_num", F.regexp_extract(F.col("Minute"), r"\d+", 0).cast("int")) # 标记是否为Minute=1的目标行 df_clean = df_clean.withColumn("is_target", F.when(F.col("Minute_num") == 1, 1).otherwise(0))
3. 生成连续序列的分组标识
通过窗口函数,按Id分区、identifier排序,对非目标行的计数进行累加,以此作为连续目标序列的分组ID:
# 定义窗口规则:按Id分区,按identifier排序保证行顺序 window = Window.partitionBy("Id").orderBy("identifier") # 生成分组ID:每遇到一个非目标行,累加值+1,同一连续目标序列的分组ID一致 df_grouped = df_clean.withColumn( "group_id", F.sum(F.when(F.col("is_target") == 0, 1).otherwise(0)).over(window) )
4. 统计每个连续序列的长度
过滤出目标行,按Id和group_id分组统计行数,最后整理输出格式:
# 过滤有效行,分组统计连续序列长度 result = df_grouped.filter(F.col("is_target") == 1) \ .groupBy("Id", "group_id") \ .agg(F.count("*").alias("Minute")) \ .drop("group_id") \ .orderBy("Id", F.col("Minute").desc()) # 排序匹配示例结果 # 展示最终结果 result.show()
最终输出
+---+------+ | Id|Minute| +---+------+ | aa| 3| | aa| 2| | bb| 2| | bb| 1| +---+------+
关键逻辑说明
- 列清洗:用
regexp_extract提取Minute中的数字,解决原始数据带备注的格式问题; - 分组标识:利用窗口累加非目标行的计数,将被非目标行分隔的连续目标行分配不同分组ID,实现连续序列的划分;
- 统计次数:通过分组计数得到每个连续1序列的长度,最终整理为需求格式。
内容的提问来源于stack exchange,提问作者JP7
相关产品推荐
相关产品推荐

