Spark数据集转换问询:按ID统计连续正数天数的方法
解决方案
要实现按ID统计mount字段连续大于0的天数,并区分不同连续分段,可通过Spark窗口函数和分组聚合完成,具体步骤如下:
1. 导入依赖函数
from pyspark.sql import Window from pyspark.sql.functions import col, row_number, sum, count, when
2. 确保数据按时间顺序排序
先对数据按id和dates排序,保证连续天数的统计基于正确的时间序列:
df_sorted = df.orderBy("id", "dates")
3. 标记有效记录并生成连续分组ID
通过窗口函数为每个id内的连续有效记录(mount>0)分配同一分组ID:
- 用
is_valid标记当前记录是否为有效(mount>0) - 用
group_id累积统计无效记录(mount<=0)的数量,连续的有效记录会共享同一个group_id
window_group = Window.partitionBy("id").orderBy("dates") df_with_group = df_sorted.withColumn( "is_valid", when(col("mount") > 0, 1).otherwise(0) ).withColumn( "group_id", sum(when(col("mount") <= 0, 1).otherwise(0)).over(window_group) )
4. 过滤无效记录
仅保留mount>0的有效记录:
df_valid = df_with_group.filter(col("is_valid") == 1)
5. 统计每个分段的连续天数
按id和group_id分组,统计每个连续分段的天数:
df_days = df_valid.groupBy("id", "group_id").agg(count("dates").alias("days"))
6. 生成分段标识some_id
为每个id内的不同连续分段生成唯一标识(若需全局唯一标识,可去掉partitionBy("id")):
window_some_id = Window.partitionBy("id").orderBy("group_id") df_final = df_days.withColumn("some_id", row_number().over(window_some_id)).select("some_id", "id", "days")
最终输出
执行上述代码后,df_final的结果与预期一致:
+-------+---+----+ |some_id|id |days| +-------+---+----+ |1 |1 |2 | |1 |2 |1 | |2 |2 |2 | |1 |3 |1 | |2 |3 |2 | +-------+---+----+
内容的提问来源于stack exchange,提问作者Camilo Romero
相关产品推荐
相关产品推荐

