如何用PySpark GroupBy+Agg计算会话的秒级时间差
PySpark计算会话时长(最新与最早时间差)
实现步骤
- 转换时间列类型:原始数据中的
timestamp是字符串格式,需先转为PySpark的TimestampType,才能进行时间运算。 - 分组聚合计算时长:按
session_id分组,计算每组内最新时间与最早时间的差值,再转换为秒数。
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql.functions import col, max, min, unix_timestamp, seconds # 初始化SparkSession spark = SparkSession.builder.appName("SessionDuration").getOrCreate() # 修正原始数据格式(移除多余嵌套括号) data = [ ("ID1", "2021-12-10 10:00:00"), ("ID1", "2021-12-10 10:05:00"), ("ID2", "2021-12-10 10:20:00"), ("ID2", "2021-12-10 10:24:00"), ("ID2", "2021-12-10 10:26:00"), ] # 创建DataFrame并转换时间列类型 df = spark.createDataFrame(data, ["session_id", "timestamp"]) df = df.withColumn("timestamp", col("timestamp").cast("timestamp")) # 计算会话时长(两种方法任选其一) # 方法1:用unix_timestamp将时间转成秒数后相减 duration_df = df.groupBy("session_id").agg( (unix_timestamp(max("timestamp")) - unix_timestamp(min("timestamp"))).alias("duration") ) # 方法2:先计算时间差(Interval类型),再提取秒数 # duration_df = df.groupBy("session_id").agg( # seconds(max("timestamp") - min("timestamp")).alias("duration") # ) # 展示结果 duration_df.show(truncate=False)
结果说明
执行代码后输出结果如下:
+----------+--------+ |session_id|duration| +----------+--------+ |ID1 |300 | |ID2 |360 | +----------+--------+
内容的提问来源于stack exchange,提问作者user2913139
相关产品推荐
相关产品推荐

