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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 14:52:47