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

PySpark按指定字段分组计算时长,首条记录标记为first

PySpark分组计算时间差并标记首条记录

需求说明

按date_id、subs_no、year、month、day分组:

  • 组内第一条记录的duration字段显示first
  • 其余记录计算当前记录与同组前一条记录的时间差,格式为时:分:秒

原始数据集

+--------+---------------+--------+----+-----+---+
| date_id|             ts| subs_no|year|month|day|
+--------+---------------+--------+----+-----+---+
|20200801|14:27:18.000000|10007239|2022|    6|  1|
|20200801|14:29:44.000000|10054647|2022|    6|  1|
|20200801|08:24:21.000000|10057750|2022|    6|  1|
|20200801|13:49:27.000000|10019958|2022|    6|  1|
|20200801|20:07:32.000000|10019958|2022|    6|  1|
+--------+---------------+--------+----+-----+---+

注:ts字段为字符串类型

预期输出

+--------+---------------+--------+----+-----+---+---------+
| date_id|             ts| subs_no|year|month|day| duration|
+--------+---------------+--------+----+-----+---+---------+
|20200801|14:27:18.000000|10007239|2022|    6|  1| first   |
|20200801|14:29:44.000000|10054647|2022|    6|  1| first   |
|20200801|08:24:21.000000|10057750|2022|    6|  1| first   |
|20200801|13:49:27.000000|10019958|2022|    6|  1| first   |
|20200801|20:07:32.000000|10019958|2022|    6|  1| 6:18:05 |
+--------+---------------+--------+----+-----+---+---------+

解决方案代码

from pyspark.sql import SparkSession
from pyspark.sql import Window
from pyspark.sql.functions import col, to_timestamp, lag, when, floor, concat, lit

# 初始化SparkSession
spark = SparkSession.builder.appName("TimeDurationCalculation").getOrCreate()

# 创建示例数据集
data = [
    ("20200801", "14:27:18.000000", "10007239", 2022, 6, 1),
    ("20200801", "14:29:44.000000", "10054647", 2022, 6, 1),
    ("20200801", "08:24:21.000000", "10057750", 2022, 6, 1),
    ("20200801", "13:49:27.000000", "10019958", 2022, 6, 1),
    ("20200801", "20:07:32.000000", "10019958", 2022, 6, 1)
]
columns = ["date_id", "ts", "subs_no", "year", "month", "day"]
df = spark.createDataFrame(data, columns)

# 1. 将字符串类型的ts转换为时间戳类型
df = df.withColumn("ts_timestamp", to_timestamp(col("ts"), "HH:mm:ss.SSSSSS"))

# 2. 定义窗口:按指定字段分区,按时间戳排序
window_spec = Window.partitionBy("date_id", "subs_no", "year", "month", "day").orderBy("ts_timestamp")

# 3. 获取前一条记录的时间戳
df = df.withColumn("prev_ts", lag("ts_timestamp", 1).over(window_spec))

# 4. 计算时间差(秒),并格式化为时:分:秒
df = df.withColumn(
    "duration",
    when(
        col("prev_ts").isNull(),  # 首条记录无前置时间
        lit("first")
    ).otherwise(
        concat(
            floor((col("ts_timestamp").cast("long") - col("prev_ts").cast("long")) / 3600).cast("string"),
            lit(":"),
            floor(((col("ts_timestamp").cast("long") - col("prev_ts").cast("long")) % 3600) / 60).cast("string"),
            lit(":"),
            ((col("ts_timestamp").cast("long") - col("prev_ts").cast("long")) % 60).cast("string")
        )
    )
)

# 5. 移除中间字段,保留目标列
result_df = df.select("date_id", "ts", "subs_no", "year", "month", "day", "duration")

# 展示结果
result_df.show()

代码说明

  • 时间戳转换:使用to_timestamp将字符串ts转为时间戳类型,方便后续计算时间差
  • 窗口定义:通过Window.partitionBy指定分组字段,orderBy确保组内记录按时间顺序排列,这样lag才能正确获取前一条记录
  • 时间差计算:将时间戳转为长整型(秒数)做差值,再通过数学运算拆分出时、分、秒,最后拼接成指定格式
  • 首条判断:用when函数判断prev_ts是否为空(即组内第一条),为空则显示first,否则显示计算出的时间差

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 13:05:24