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

如何筛选包含至少3个5分钟间隔的Spark窗口分区

筛选包含至少3个5分钟间隔的Spark窗口分区

需求概述

筛选出按subs_no、year、month、day、cgi分组后,分组内存在至少3条连续记录且相邻间隔≤5分钟的分区,保留该分区的所有记录。

现有代码

from pyspark.sql import functions as F, Window as W

df_subs_loc_movmnt_ts = df_subs_loc_movmnt.withColumn("new_ts", F.unix_timestamp(F.col("ts"), "HH:mm:ss"))
w = W.partitionBy('subs_no', 'year', 'month', 'day', 'cgi').orderBy('new_ts')
df_subs_loc_movmnt_duration = df_subs_loc_movmnt_ts.withColumn('duration', F.from_unixtime(F.col('new_ts') - F.min('new_ts').over(w), "HH:mm:ss"))

原始DataFrame输出

+--------+---------------+--------+---------------+-------------+----+-----+---+------+--------+
| date_id|             ts| subs_no|            cgi|       msisdn|year|month|day|new_ts|duration|
+--------+---------------+--------+---------------+-------------+----+-----+---+------+--------+
|20200801|10:40:43.000000|10000093|510-11-610725-5|7664622154085|2022|    6|  2| 13243|07:00:00|
|20200801|12:55:30.000000|10000093|510-11-610725-5|7664622154085|2022|    6|  2| 21330|09:14:47|
|20200801|05:30:47.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| -5353|07:00:00|
|20200801|10:55:21.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 14121|12:24:34|
|20200801|13:05:06.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 21906|14:34:19|
|20200801|13:05:50.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 21950|14:35:03|
|20200801|13:06:49.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 22009|14:36:02|
|20200801|13:08:32.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 22112|14:37:45|
|20200801|13:08:44.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 22124|14:37:57|
|20200801|13:09:01.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 22141|14:38:14|
|20200801|19:09:51.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 43791|20:39:04|
|20200801|19:37:16.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 45436|21:06:29|
|20200801|19:55:17.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 46517|21:24:30|
|20200801|13:24:58.000000|10000393|510-11-610354-1|4745471710184|2022|    6|  2| 23098|07:00:00|
|20200801|13:54:43.000000|10000393|510-11-610372-4|4745471710184|2022|    6|  3| 24883|07:00:00|
|20200801|11:38:41.000000|10000406|510-11-620412-4|4411875524723|2022|    6|  3| 16721|07:00:00|
|20200801|08:38:36.000000|10000514|510-11-610658-6|5908140424233|2022|    6|  2|  5916|07:00:00|
|20200801|02:12:05.000000|10000610|510-11-610030-9|6354719688724|2022|    6|  1|-17275|07:00:00|
|20200801|06:41:58.000000|10000610|510-11-610030-9|6354719688724|2022|    6|  1| -1082|11:29:53|
|20200801|06:51:14.000000|10000610|510-11-610030-9|6354719688724|2022|    6|  1|  -526|11:39:09|
+--------+---------------+--------+---------------+-------------+----+-----+---+------+--------+

期望筛选结果

+--------+---------------+--------+---------------+-------------+----+-----+---+------+--------+
| date_id|             ts| subs_no|            cgi|       msisdn|year|month|day|new_ts|duration|
+--------+---------------+--------+---------------+-------------+----+-----+---+------+--------+
|20200801|05:30:47.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| -5353|07:00:00|
|20200801|10:55:21.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 14121|12:24:34|
|20200801|13:05:06.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 21906|14:34:19|
|20200801|13:05:50.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 21950|14:35:03|
|20200801|13:06:49.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 22009|14:36:02|
|20200801|13:08:32.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 22112|14:37:45|
|20200801|13:08:44.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 22124|14:37:57|
|20200801|13:09:01.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 22141|14:38:14|
|20200801|19:09:51.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 43791|20:39:04|
|20200801|19:37:16.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 45436|21:06:29|
|20200801|19:55:17.000000|10000118|510-11-610195-5|7560242795888|2022|    6|  2| 46517|21:24:30|
+--------+---------------+--------+---------------+-------------+----+-----+---+------+--------+

DataFrame Schema

StructType(List(
    StructField(date_id, IntegerType, true),
    StructField(ts, StringType, true),
    StructField(subs_no, LongType, true),
    StructField(cgi, StringType, true),
    StructField(msisdn, LongType, true),
    StructField(year, IntegerType, true),
    StructField(month, IntegerType, true),
    StructField(day, IntegerType, true),
    StructField(new_ts, LongType, true),
    StructField(duration, StringType, true)
))

实现方案

步骤说明

  1. 计算分组内每条记录与前一条记录的时间间隔(秒)。
  2. 统计每个分组内,间隔≤5分钟(300秒)的次数。
  3. 筛选出间隔达标(至少2次,对应3条连续记录在10分钟内,覆盖3个5分钟区间)的分组,保留其所有记录。

完整代码

from pyspark.sql import functions as F, Window as W

# 现有代码生成df_subs_loc_movmnt_duration
df_subs_loc_movmnt_ts = df_subs_loc_movmnt.withColumn("new_ts", F.unix_timestamp(F.col("ts"), "HH:mm:ss"))
w_part_order = W.partitionBy('subs_no', 'year', 'month', 'day', 'cgi').orderBy('new_ts')
df_subs_loc_movmnt_duration = df_subs_loc_movmnt_ts.withColumn('duration', F.from_unixtime(F.col('new_ts') - F.min('new_ts').over(w_part_order), "HH:mm:ss"))

# 计算相邻记录的时间间隔
df_with_interval = df_subs_loc_movmnt_duration.withColumn(
    'prev_new_ts',
    F.lag('new_ts').over(w_part_order)
).withColumn(
    'interval',
    F.when(F.col('prev_new_ts').isNotNull(), F.col('new_ts') - F.col('prev_new_ts')).otherwise(0)
)

# 统计分组内有效间隔(≤300秒)的数量
w_group = W.partitionBy('subs_no', 'year', 'month', 'day', 'cgi')
df_with_stats = df_with_interval.withColumn(
    'valid_interval_count',
    F.sum(F.when(F.col('interval') <= 300, 1).otherwise(0)).over(w_group)
)

# 筛选出符合条件的分组(至少2个有效间隔,对应至少3条连续记录在10分钟内)
result_df = df_with_stats.filter(F.col('valid_interval_count') >= 2).drop('prev_new_ts', 'interval', 'valid_interval_count')

# 查看结果
result_df.show()

代码说明

  • lag('new_ts').over(w_part_order):获取分组内前一条记录的时间戳,用于计算间隔。
  • valid_interval_count:统计分组内相邻时间间隔≤5分钟的次数,次数≥2意味着存在至少3条连续记录的间隔都在5分钟内,满足“包含至少3个5分钟间隔”的需求。
  • 最后
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 23:45:39