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

如何在PySpark中基于时间戳和userId生成sessionId会话标识列

实现思路

  • 按用户ID分组,组内按时间戳升序排序,计算每条记录与同用户上一条记录的时间差
  • 若记录是用户的第一条行为,或与上一条行为的时间差超过30分钟(1800*1000毫秒),标记为新会话起点
  • 对新会话标记累加得到用户内部的会话序号,再按会话起始时间全局排序生成全局统一的sessionId

完整代码

from pyspark.sql import functions as F 
from pyspark.sql.window import Window

# 示例数据
df = spark.createDataFrame([
  ("blue", "view", 1610494094750, 11),
  ("green", "add to bag", 1510593114350, 21),
  ("red", "close", 1610493115350, 41),
  ("blue", "view", 1610494094350, 11),
  ("blue", "close", 1510593114312, 21),
  ("red", "view", 1610493114350, 41),
  ("red", "view", 1610593114350, 41),
  ("green", "purchase", 1610494094350, 31)
], ["item", "event", "timestamp", "userId"])

# 步骤1:按用户分区、时间排序的窗口
w_user_time = Window.partitionBy("userId").orderBy("timestamp")

# 步骤2:获取同用户上一条行为的时间戳
df = df.withColumn("prev_timestamp", F.lag("timestamp").over(w_user_time))

# 步骤3:标记新会话起点
df = df.withColumn("is_new_session", F.when(
    F.col("prev_timestamp").isNull() | (F.col("timestamp") - F.col("prev_timestamp") > 1800 * 1000),
    1
).otherwise(0))

# 步骤4:计算用户内部的会话序号
df = df.withColumn("user_session_seq", F.sum("is_new_session").over(w_user_time))

# 步骤5:获取每个会话的最早时间,用于全局排序生成统一sessionId
w_session = Window.partitionBy("userId", "user_session_seq")
df = df.withColumn("session_start_ts", F.min("timestamp").over(w_session))

# 步骤6:全局按会话起始时间排序,生成全局session序号
w_global_session = Window.orderBy("session_start_ts", "userId", "user_session_seq")
df = df.withColumn("global_session_seq", F.dense_rank().over(w_global_session))

# 步骤7:生成最终sessionId
df = df.withColumn("sessionId", F.concat(F.lit("session"), F.col("global_session_seq")))

# 输出结果,可按需要排序
result_df = df.select("item", "event", "timestamp", "userId", "sessionId").orderBy("global_session_seq", "timestamp")
result_df.show(truncate=False)

输出结果与你给出的预期完全一致:

+-----+----------+-------------+------+---------+
|item |event     |timestamp    |userId|sessionId|
+-----+----------+-------------+------+---------+
|blue |close     |1510593114312|21    |session1 |
|green|add to bag|1510593114350|21    |session1 |
|red  |view      |1610493114350|41    |session2 |
|red  |close     |1610493115350|41    |session2 |
|blue |view      |1610494094350|11    |session3 |
|blue |view      |1610494094750|11    |session3 |
|green|purchase  |1610494094350|31    |session4 |
|red  |view      |1610593114350|41    |session5 |
+-----+----------+-------------+------+---------+

性能优化说明

如果是大数据量场景,不需要全局统一的sessionX格式序号,建议直接使用concat(F.col("userId"), F.lit("_"), F.col("user_session_seq"))作为sessionId,可避免全局排序带来的性能损耗。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 10:27:03