PySpark按客户分组填充递增day_number的多方案问询
PySpark填充day_number空值的两种实现方案(row_number/last函数)
假设我们有如下示例DataFrame,每个客户的日期是连续的,day_number字段存在部分非空值,其余为NULL,需要将NULL填充为按日期递增1的序列值:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window spark = SparkSession.builder.appName("FillDayNumber").getOrCreate() data = [ ("customer1", "2023-05-13", "x1", 1), ("customer1", "2023-05-14", "x2", None), ("customer1", "2023-05-15", "x3", None), ("customer1", "2023-05-16", "x4", 4), ("customer1", "2023-05-17", "x5", None), ("customer2", "2023-06-01", "y1", None), ("customer2", "2023-06-02", "y2", 2), ("customer2", "2023-06-03", "y3", None), ("customer2", "2023-06-04", "y4", None), ] df = spark.createDataFrame(data, ["customer", "date", "col_x", "day_number"]) df.show()
方案一:基于row_number的锚点递推法
核心思路是用row_number()标记每个客户内的行顺序,结合last()函数捕捉最近的非空day_number作为锚点,通过行号差计算递推值。如果客户没有任何非空的day_number,直接用row_number作为序列值。
# 1. 按客户分组、日期排序,添加行号列 window = Window.partitionBy("customer").orderBy("date") df_with_rn = df.withColumn("rn", F.row_number().over(window)) # 2. 向前传播最近的非空day_number及其对应的行号 window_anchor = Window.partitionBy("customer").orderBy("date").rowsBetween(Window.unboundedPreceding, Window.currentRow) df_with_anchor = df_with_rn.withColumn("anchor_day", F.last("day_number", ignorenull=True).over(window_anchor)) \ .withColumn("anchor_rn", F.last(F.when(F.col("day_number").isNotNull(), F.col("rn")), ignorenull=True).over(window_anchor)) # 3. 计算填充后的day_number df_filled = df_with_anchor.withColumn("filled_day_number", F.when(F.col("anchor_day").isNotNull(), F.col("anchor_day") + (F.col("rn") - F.col("anchor_rn"))) .otherwise(F.col("rn"))) \ .drop("rn", "anchor_day", "anchor_rn") df_filled.show()
方案二:基于日期差的last函数填充法
利用题目中“每个客户日期连续”的前提,通过计算当前日期与锚点日期的天数差,结合最近的非空day_number递推。无锚点时,以客户的第一个日期为起点计算序列。
# 1. 将字符串日期转换为日期类型 df_date = df.withColumn("date", F.to_date("date")) # 2. 向前传播最近的非空day_number及其对应日期 window = Window.partitionBy("customer").orderBy("date") window_anchor = Window.partitionBy("customer").orderBy("date").rowsBetween(Window.unboundedPreceding, Window.currentRow) df_with_anchor = df_date.withColumn("anchor_day", F.last("day_number", ignorenull=True).over(window_anchor)) \ .withColumn("anchor_date", F.last(F.when(F.col("day_number").isNotNull(), F.col("date")), ignorenull=True).over(window_anchor)) # 3. 获取每个客户的第一个日期,用于无锚点的情况 window_first = Window.partitionBy("customer") df_with_first = df_with_anchor.withColumn("first_date", F.first("date").over(window_first)) # 4. 计算填充后的day_number df_filled = df_with_first.withColumn("filled_day_number", F.when(F.col("anchor_day").isNotNull(), F.col("anchor_day") + F.datediff(F.col("date"), F.col("anchor_date"))) .otherwise(F.datediff(F.col("date"), F.col("first_date")) + 1)) \ .drop("anchor_day", "anchor_date", "first_date") df_filled.show()
两种方案的适用场景
- 方案一不依赖日期连续性,即使日期存在断档,也能保证
filled_day_number严格连续递增,适合对序列连续性要求高的场景。 - 方案二依赖日期连续的前提,计算逻辑更贴合“按日期自然递推”的业务直觉,性能略优(无需维护行号)。
内容的提问来源于stack exchange,提问作者amstergc20
相关产品推荐
相关产品推荐

