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

