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
相关产品推荐
相关产品推荐

