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

PySpark中如何按时间为数据分区分配正确的排名

问题描述

现有如下PySpark DataFrame:

simpleData = (("U1", "cd1dd155-ccd8-4b8c-bea7-571359e35fed", 1655605947), \
    ("U1", "7f20182f-8c82-4c70-8213-f7889cfdd5eb", 1655777060),  \
    ("U1", "7f20182f-8c82-4c70-8213-f7889cfdd5eb", 1655777062),
    ("U1", "c4d5a218-d61d-4e9a-b1ea-646f676c4cb7", 1656209951),   \
    ("U1", "c4d5a218-d61d-4e9a-b1ea-646f676c4cb7", 1656209952),  \
    ("U1", "c4d5a218-d61d-4e9a-b1ea-646f676c4cb7", 1656209999),    \
  )
 
columns= ["UID", "Sess", "Time"]
df = spark.createDataFrame(data = simpleData, schema = columns)
df.printSchema()
df.show(truncate=False)

输出结果:

+---+------------------------------------+----------+
|UID|Sess                                |Time      |
+---+------------------------------------+----------+
|U1 |cd1dd155-ccd8-4b8c-bea7-571359e35fed|1655605947|
|U1 |7f20182f-8c82-4c70-8213-f7889cfdd5eb|1655777060|
|U1 |7f20182f-8c82-4c70-8213-f7889cfdd5eb|1655777062|
|U1 |c4d5a218-d61d-4e9a-b1ea-646f676c4cb7|1656209951|
|U1 |c4d5a218-d61d-4e9a-b1ea-646f676c4cb7|1656209952|
|U1 |c4d5a218-d61d-4e9a-b1ea-646f676c4cb7|1656209999|
+---+------------------------------------+----------+

当前使用以下代码为窗口分区内的行分配排名:

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

df2 = df.withColumn("sess_2", F.dense_rank().over(Window.orderBy('UID', 'Sess')))
df2.show(truncate=False)

得到的输出:

+---+------------------------------------+----------+------+
|UID|Sess                                |Time      |sess_2|
+---+------------------------------------+----------+------+
|U1 |7f20182f-8c82-4c70-8213-f7889cfdd5eb|1655777060|1     |
|U1 |7f20182f-8c82-4c70-8213-f7889cfdd5eb|1655777062|1     |
|U1 |c4d5a218-d61d-4e9a-b1ea-646f676c4cb7|1656209951|2     |
|U1 |c4d5a218-d61d-4e9a-b1ea-646f676c4cb7|1656209952|2     |
|U1 |c4d5a218-d61d-4e9a-b1ea-646f676c4cb7|1656209999|2     |
|U1 |cd1dd155-ccd8-4b8c-bea7-571359e35fed|1655605947|3     |
+---+------------------------------------+----------+------+

期望输出:

+---+------------------------------------+----------+------+
|UID|Sess                                |Time      |sess_2|
+---+------------------------------------+----------+------+
|U1 |cd1dd155-ccd8-4b8c-bea7-571359e35fed|1655605947|1     |
|U1 |7f20182f-8c82-4c70-8213-f7889cfdd5eb|1655777060|2     |
|U1 |7f20182f-8c82-4c70-8213-f7889cfdd5eb|1655777062|2     |
|U1 |c4d5a218-d61d-4e9a-b1ea-646f676c4cb7|1656209951|3     |
|U1 |c4d5a218-d61d-4e9a-b1ea-646f676c4cb7|1656209952|3     |
|U1 |c4d5a218-d61d-4e9a-b1ea-646f676c4cb7|1656209999|3     |
+---+------------------------------------+----------+------+

需要实现按时间顺序为每个UID和Sess分区分配正确的sess_2排名。

解决方案

原代码的问题在于按Sess字符串的字典序排序分配排名,而非会话的实际出现时间顺序。正确逻辑是基于每个UID+Sess组合的最早时间来确定排名顺序,以下提供两种实现方式:

方式一:分组聚合+关联

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

# 1. 计算每个UID和Sess的最早出现时间
sess_min_time = df.groupBy("UID", "Sess").agg(F.min("Time").alias("min_time"))

# 2. 基于UID和最早时间对会话分配dense_rank
rank_window = Window.partitionBy("UID").orderBy("min_time")
sess_ranked = sess_min_time.withColumn("sess_2", F.dense_rank().over(rank_window))

# 3. 关联回原DataFrame得到结果
result_df = df.join(sess_ranked, on=["UID", "Sess"], how="left")
result_df.show(truncate=False)

方式二:嵌套窗口函数(无需关联)

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

# 1. 计算每个UID+Sess组合的最早时间
min_time_window = Window.partitionBy("UID", "Sess")
df_with_min_time = df.withColumn("min_time", F.min("Time").over(min_time_window))

# 2. 基于UID和最早时间分配排名,最后删除临时列
rank_window = Window.partitionBy("UID").orderBy("min_time")
result_df = df_with_min_time.withColumn("sess_2", F.dense_rank().over(rank_window)).drop("min_time")
result_df.show(truncate=False)

两种方式都能得到期望的输出,核心是用会话的最早出现时间作为排序依据,而非会话ID的字符串顺序。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 23:10:31