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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 12:05:32