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

使用PySpark窗口函数或SQL实现基于前值的时间累加计算

问题描述

现有输入数据集:

+----+----+-------------------+
|orig|dest|   actl_dep_lcl_tms|
+----+----+-------------------+
| YYZ| YYR|2022-12-31 08:50:00|
| YYZ| YYR|2022-12-31 10:18:00|
| YYZ| YYR|2022-12-31 11:23:00|
+----+----+-------------------+

需要生成包含Output_col的输出数据集,规则为:

  • Output_col的起始值是orig+dest分组内actl_dep_lcl_tms的最小值
  • 后续每行的Output_col = 前一行Output_col值 + 15分钟

注:原示例中的Output_col数值不符合累加规则,推测为笔误,以下解法严格按需求实现。

解决方案

方法1:Spark SQL

通过窗口函数分组获取基准时间,再按行号累加时间间隔:

WITH grouped_data AS (
    SELECT 
        orig,
        dest,
        actl_dep_lcl_tms,
        -- 获取分组内最小时间作为基准
        MIN(actl_dep_lcl_tms) OVER (PARTITION BY orig, dest) AS base_time,
        -- 给分组内每行编号
        ROW_NUMBER() OVER (PARTITION BY orig, dest ORDER BY actl_dep_lcl_tms) AS rn
    FROM your_table
)
SELECT 
    orig,
    dest,
    actl_dep_lcl_tms,
    -- 按行号累加15分钟
    DATE_ADD(base_time, INTERVAL (rn - 1)*15 MINUTE) AS Output_col
FROM grouped_data
ORDER BY orig, dest, actl_dep_lcl_tms;

方法2:PySpark DataFrame API

用DataFrame链式调用实现相同逻辑:

from pyspark.sql import Window
from pyspark.sql.functions import min, row_number, expr

# 定义窗口
group_window = Window.partitionBy("orig", "dest")
order_window = Window.partitionBy("orig", "dest").orderBy("actl_dep_lcl_tms")

# 计算基准时间、行号,生成Output_col
result_df = df.withColumn("base_time", min("actl_dep_lcl_tms").over(group_window)) \
              .withColumn("rn", row_number().over(order_window)) \
              .withColumn("Output_col", expr("date_add(base_time, interval (rn-1)*15 minute)")) \
              .select("orig", "dest", "actl_dep_lcl_tms", "Output_col")

方法3:Pandas

通过分组变换和时间增量实现:

import pandas as pd

# 构造数据集(实际可从文件/数据库读取)
df = pd.DataFrame({
    "orig": ["YYZ", "YYZ", "YYZ"],
    "dest": ["YYR", "YYR", "YYR"],
    "actl_dep_lcl_tms": pd.to_datetime(["2022-12-31 08:50:00", "2022-12-31 10:18:00", "2022-12-31 11:23:00"])
})

# 计算分组基准时间、行号,生成Output_col
df["base_time"] = df.groupby(["orig", "dest"])["actl_dep_lcl_tms"].transform("min")
df["rn"] = df.groupby(["orig", "dest"]).cumcount()
df["Output_col"] = df["base_time"] + pd.to_timedelta(df["rn"] * 15, unit="m")

# 保留目标列
result_df = df[["orig", "dest", "actl_dep_lcl_tms", "Output_col"]]

内容的提问来源于stack exchange,提问作者user1920851

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 09:37:17