PySpark实现24小时窗口内重复状态消息去重方案咨询
精简PySpark DataFrame:基于状态的非连续24小时窗口去重
需求描述
现有包含status(整数类型)和ts(时间戳类型)的PySpark DataFrame,需移除同一状态在某个「新」窗口的24小时内重复出现的行,规则如下:
- 某状态的首个24小时窗口从该状态的第一条消息的时间戳开始
- 该状态的下一个24小时窗口,从首个窗口结束时间后的第一条该状态消息的时间戳开始(窗口非连续)
示例数据
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, IntegerType, TimestampType import datetime spark = SparkSession.builder.appName("StatusDeduplication").getOrCreate() data = [(10, datetime.datetime.strptime("2022-01-01 00:00:00", "%Y-%m-%d %H:%M:%S")), (10, datetime.datetime.strptime("2022-01-01 04:00:00", "%Y-%m-%d %H:%M:%S")), (10, datetime.datetime.strptime("2022-01-01 23:00:00", "%Y-%m-%d %H:%M:%S")), (10, datetime.datetime.strptime("2022-01-02 05:00:00", "%Y-%m-%d %H:%M:%S")), (10, datetime.datetime.strptime("2022-01-02 06:00:00", "%Y-%m-%d %H:%M:%S")), (20, datetime.datetime.strptime("2022-01-01 03:00:00", "%Y-%m-%d %H:%M:%S")) ] myschema = StructType( [ StructField("status", IntegerType()), StructField("ts", TimestampType()) ] ) df = spark.createDataFrame(data=data, schema=myschema)
窗口规则说明
- status=10的首个窗口:2022-01-01 00:00:00 ~ 2022-01-02 00:00:00
- status=10的第二个窗口:2022-01-02 05:00:00 ~ 2022-01-03 05:00:00
- status=20的首个窗口:2022-01-01 03:00:00 ~ 2022-01-02 03:00:00
期望结果
expected_data = [(10, datetime.datetime.strptime("2022-01-01 00:00:00", "%Y-%m-%d %H:%M:%S")), (10, datetime.datetime.strptime("2022-01-02 05:00:00", "%Y-%m-%d %H:%M:%S")), (20, datetime.datetime.strptime("2022-01-01 03:00:00", "%Y-%m-%d %H:%M:%S")) ] expected_df = spark.createDataFrame(data=expected_data, schema=myschema)
Window函数实现方案
核心思路是为每个状态的记录划分非连续的24小时窗口组,再保留每组的第一条记录。具体步骤如下:
步骤1:按状态分区、时间排序,添加行号
先对每个状态下的记录按时间戳排序,标记行号方便后续计算:
from pyspark.sql.window import Window from pyspark.sql import functions as F window_partition = Window.partitionBy("status").orderBy("ts") df_with_row = df.withColumn("row_num", F.row_number().over(window_partition))
步骤2:计算每个记录所属的窗口起始时间
使用F.last函数结合条件判断,递归确定每条记录的窗口起始时间:
- 第一条记录的窗口起始为自身的
ts - 后续记录若
ts>= 上一个窗口起始 + 24小时,则以自身ts作为新窗口起始;否则继承上一个窗口的起始
df_with_window_start = df_with_row.withColumn( "window_start", F.when( F.col("row_num") == 1, F.col("ts") ).otherwise( F.when( F.col("ts") >= F.expr("last(window_start) over (partition by status order by ts) + interval 24 hours"), F.col("ts") ).otherwise( F.expr("last(window_start) over (partition by status order by ts)") ) ) )
步骤3:按状态和窗口起始分组,保留每组第一条记录
每个窗口组只保留起始时间对应的那条记录,完成去重:
deduplicated_df = df_with_window_start.groupBy("status", "window_start")\ .agg(F.first("ts").alias("ts"))\ .select("status", "ts")\ .orderBy("status", "ts")
运行后deduplicated_df的输出与期望结果一致。
内容的提问来源于stack exchange,提问作者Cribber
相关产品推荐
相关产品推荐

