在PySpark中获取CPO_BB_Status=1的时间窗口及时长
PySpark/Databricks 状态时间窗口计算实现
需求说明
按CpoSku分组,从每个SKU的最早DateUpdated记录开始,识别CPO_BB_Status连续为1的时间窗口,计算每个窗口从状态变为1到转为0的时间间隔,输出窗口的起始时间、结束时间、秒级时长和分钟级时长。
数据重建代码
首先执行以下代码重建测试数据集(已修正原数据中的时区笔误):
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, TimestampType, IntegerType from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession(Databricks环境可省略此步骤,直接使用spark对象) spark = SparkSession.builder \ .appName("StatusWindowCalculation") \ .getOrCreate() # 定义Schema schema = StructType([ StructField("CpoSku", StringType(), True), StructField("DateUpdated", TimestampType(), True), StructField("CPO_BB_Status", IntegerType(), True) ]) # 测试数据(修正了时区笔误:+00:01改为+00:00) data = [ ("AAGN7013005", "2024-01-24T05:02:06.898+00:00", 0), ("AAGN7013005", "2024-01-24T05:07:05.090+00:00", 1), ("AAGN7013005", "2024-01-24T06:42:56.825+00:00", 1), ("AAGN7013005", "2024-01-24T06:48:01.647+00:00", 1), ("AAGN7013005", "2024-01-24T07:48:18.456+00:00", 1), ("AAGN7013005", "2024-01-24T09:30:22.534+00:00", 1), ("AAGN7013005", "2024-01-24T09:36:04.075+00:00", 1), ("AAGN7013005", "2024-01-24T10:39:04.796+00:00", 1), ("AAGN7013005", "2024-01-24T10:44:01.193+00:00", 1), ("AAGN7013005", "2024-01-24T17:00:06.217+00:00", 1), ("AAGN7013005", "2024-01-24T18:07:16.612+00:00", 1), ("AAGN7013005", "2024-01-24T18:13:04.639+00:00", 0), ("AAGN7013005", "2024-01-24T21:33:03.796+00:00", 0), ("AAGN7013005", "2024-01-24T21:38:28.834+00:00", 1), ("AAGN7013005", "2024-01-24T22:35:43.995+00:00", 1), ("AAGN7013005", "2024-01-24T22:40:45.930+00:00", 0), ("AAGN7022205", "2024-01-24T04:09:30.167+00:00", 0), ("AAGN7022205", "2024-01-24T04:14:56.294+00:00", 0), ("AAGN7022205", "2024-01-24T04:53:01.281+00:00", 0), ("AAGN7022205", "2024-01-24T05:03:27.103+00:00", 0), ("AAGN7022205", "2024-01-24T05:08:05.096+00:00" ,1), ("AAGN7022205", "2024-01-24T05:53:22.652+00:00", 1), ("AAGN7022205", "2024-01-24T06:04:59.031+00:00", 1), ("AAGN7022205", "2024-01-24T06:43:04.285+00:00", 1), ("AAGN7022205", "2024-01-24T06:43:34.285+00:00", 0) ] # 创建DataFrame df_test = spark.createDataFrame(data, schema=schema)
核心实现代码
以下是计算状态窗口的完整逻辑:
# 1. 定义窗口:按CpoSku分组,按DateUpdated排序 window_spec = Window.partitionBy("CpoSku").orderBy("DateUpdated") # 2. 标记状态切换点:当前状态与前一条不同时标记为1 df_with_switch = df_test.withColumn( "status_switch", F.when(F.col("CPO_BB_Status") != F.lag("CPO_BB_Status").over(window_spec), 1).otherwise(0) ) # 3. 为每个连续状态窗口分配唯一ID df_with_window_id = df_with_switch.withColumn( "window_id", F.sum("status_switch").over(window_spec.rangeBetween(Window.unboundedPreceding, 0)) ) # 4. 按CpoSku和window_id分组,计算每个窗口的起止时间和状态 window_summary = df_with_window_id.groupBy("CpoSku", "window_id", "CPO_BB_Status") \ .agg( F.min("DateUpdated").alias("window_start"), F.max("DateUpdated").alias("window_end") ) # 5. 筛选出状态为1的窗口,计算时长 result_df = window_summary.filter(F.col("CPO_BB_Status") == 1) \ .withColumn("duration_seconds", F.unix_timestamp("window_end") - F.unix_timestamp("window_start")) \ .withColumn("duration_minutes", F.round(F.col("duration_seconds") / 60, 2)) \ .select("CpoSku", "window_start", "window_end", "duration_seconds", "duration_minutes") # 展示结果 result_df.orderBy("CpoSku", "window_start").show(truncate=False)
代码解释
- 步骤1-2:通过
lag函数对比当前与前一条记录的状态,标记状态切换的位置。 - 步骤3:对切换标记进行累加,为每个连续的相同状态生成唯一的
window_id,实现状态窗口的划分。 - 步骤4:按SKU和窗口ID分组,提取每个窗口的最早(起始)和最晚(结束)时间。
- 步骤5:筛选出状态为1的窗口,用
unix_timestamp计算秒级时长,再转换为分钟级时长。
预期输出
+-----------+-----------------------+-----------------------+---------------+----------------+ |CpoSku |window_start |window_end |duration_seconds|duration_minutes| +-----------+-----------------------+-----------------------+---------------+----------------+ |AAGN7013005|2024-01-24 05:07:05.09 |2024-01-24 18:13:04.639|47159 |785.98 | |AAGN7013005|2024-01-24 21:38:28.834|2024-01-24 22:40:45.93 |3737 |62.28 | |AAGN7022205|2024-01-24 05:08:05.096|2024-01-24 06:43:34.285|5729 |95.48 | +-----------+-----------------------+-----------------------+---------------+----------------+
内容的提问来源于stack exchange,提问作者irum zahra
相关产品推荐
相关产品推荐

