如何使用Spark统计指定历史时间段内的重复值数量
24小时窗口重复值统计Spark实现方案
核心逻辑
用Spark的时间范围窗口函数实现:按目标编码列分区,按时间升序排序,窗口范围限定为当前行往前24小时到当前行,统计窗口内的总行数后减1(排除当前行本身),即可得到对应时间区间内的重复值数量。
完整实现代码
首先导入依赖、构造测试数据集:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义列名 columns = ["mcc", "event_time"] # 测试数据集 data = [ ("5812", "2020-12-27T17:28:32.000+0000"), ("5812", "2020-12-25T17:35:32.000+0000"), ("5812", "2020-12-25T13:04:05.000+0000"), ("7999", "2020-12-25T09:23:01.000+0000"), ("5999","2020-12-25T07:29:52.000+0000"), ("5814", "2020-12-25T12:23:05.000+0000"), ("5814", "2020-12-25T11:52:57.000+0000"), ("5814", "2020-12-24T11:00:57.000+0000"), ("5999", "2020-12-24T07:29:52.000+0000") ] df = spark.createDataFrame(data).toDF(*columns)
然后做时间类型转换、定义窗口计算重复值:
# 字符串时间转为timestamp类型 df = df.withColumn("event_ts", F.to_timestamp("event_time")) # 定义范围窗口:按mcc分区,按秒级时间戳排序,窗口覆盖当前行往前24小时到当前行 time_window = Window.partitionBy("mcc") \ .orderBy(F.unix_timestamp("event_ts")) \ .rangeBetween(-86400, Window.currentRow) # 统计重复次数:窗口内总行数减1(排除当前行本身) df_result = df.withColumn("repeat_count_24h", F.count("*").over(time_window) - 1) \ .orderBy("mcc", "event_ts") # 查看结果 df_result.show(truncate=False)
结果说明
执行后输出的repeat_count_24h列就是每条记录对应过去24小时内同mcc的重复出现次数,和预期输出完全匹配。如果需要统计全量重复值总个数,直接过滤repeat_count_24h >= 1的行后按mcc去重计数即可。
内容的提问来源于stack exchange,提问作者user1389739
相关产品推荐
相关产品推荐

