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

如何用PySpark GroupBy+Agg计算会话的秒级时间差

PySpark计算会话时长(最新与最早时间差)

实现步骤

  1. 转换时间列类型:原始数据中的timestamp是字符串格式,需先转为PySpark的TimestampType,才能进行时间运算。
  2. 分组聚合计算时长:按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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 02:01:57