如何筛选包含至少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) ))
实现方案
步骤说明
- 计算分组内每条记录与前一条记录的时间间隔(秒)。
- 统计每个分组内,间隔≤5分钟(300秒)的次数。
- 筛选出间隔达标(至少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分钟间隔”的需求。- 最后
相关产品推荐
相关产品推荐

