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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 21:27:32