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

PySpark按Id分组统计连续出现的Minute=1的次数

PySpark实现连续Minute=1序列的次数统计

问题描述

给定如下DataFrame:

IdMinuteidentifier
aa11
aa12
aa13
aa2(忽略)4
aa15
aa16
bb17
bb18
bb5(忽略)9
bb110

需求:按Id分组,统计连续出现的Minute值为1的序列长度,输出每个连续序列对应的Id和次数,期望结果如下:

IdMinute
aa3
aa2
bb2
bb1

实现步骤与代码

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 10:25:25