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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 10:45:48