如何在PySpark中为客户及订单相关字段生成分组行号
问题分析与解决方案
你的代码不符合预期的核心原因是窗口分区规则错误:你将customer_id、order_ts、order_nbr三个字段作为分区键,这会把每个独立的订单组(同一用户、同一订单时间、同一订单号)拆分为单独分区,row_number()会在每个分区内为每行生成递增序号(1、2、3...),导致同一订单组内的row_num无法保持一致。
正确实现代码
要实现同一订单组内row_num值相同,且按订单时间从新到旧为每个用户的订单组连续编号,需调整窗口规则并改用dense_rank()函数:
from pyspark.sql import Window from pyspark.sql.functions import col, dense_rank # 定义窗口:按用户分区,按订单时间倒序排序 window_spec = Window.partitionBy("customer_id").orderBy(col("order_ts").desc()) # 添加row_num列 df1 = df.withColumn("row_num", dense_rank().over(window_spec))
原理说明
partitionBy("customer_id"):将数据按用户分组,确保编号仅在同一用户范围内生效。orderBy(col("order_ts").desc()):对每个用户的订单按时间从新到旧排序,最新的订单组会获得最小的编号。dense_rank():为同一用户下的每个订单组(相同order_ts和order_nbr)分配相同的连续编号,不会出现跳号,完全匹配需求效果。
内容的提问来源于stack exchange,提问作者user3642360
相关产品推荐
相关产品推荐

