如何在PySpark中为缺失的date_id补全对应行数据
PySpark 补全缺失日期并继承最近前置数据
核心思路
要实现缺失日期的补全并继承最近前置日期的对应数据,可通过以下步骤完成:
- 生成覆盖数据中最小到最大日期的完整日期序列
- 为每个唯一的
anchor生成与完整日期序列的笛卡尔积,确保每个anchor在所有日期都有对应记录 - 利用窗口函数或日期区间生成的方式,为缺失日期填充最近的前置有效数据
方法一:窗口函数向前填充(逻辑直观,适合小数据量)
1. 初始化测试数据
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window spark = SparkSession.builder.appName("FillMissingDates").getOrCreate() # 构造测试数据集 data = [ ("2023-10-01", 123, ["456","345"]), ("2023-10-01", 256, ["278","875"]), ("2023-10-03", 123, ["703","409"]), ("2023-10-03", 256, ["801","704"]), ("2023-10-06", 123, ["398","209"]), ("2023-10-06", 256, ["659","783"]) ] df = spark.createDataFrame(data, ["date_id", "anchor", "subs_ids"]) # 将字符串类型的date_id转为日期类型 df = df.withColumn("date_id", F.to_date("date_id"))
2. 生成完整日期序列
# 获取数据中的最小和最大日期 date_bounds = df.select(F.min("date_id").alias("min_date"), F.max("date_id").alias("max_date")).collect()[0] min_date, max_date = date_bounds["min_date"], date_bounds["max_date"] # 生成从min_date到max_date的所有连续日期 full_dates = spark.sql(f""" SELECT sequence(to_date('{min_date}'), to_date('{max_date}'), interval 1 day) AS date_array """).select(F.explode("date_array").alias("date_id"))
3. 生成所有(anchor, date_id)组合并填充数据
# 获取所有唯一的anchor值 unique_anchors = df.select("anchor").distinct() # 生成anchor与完整日期的笛卡尔积,得到所有可能的组合 full_combinations = unique_anchors.crossJoin(full_dates) # 关联原始数据,标记有效记录 combined = full_combinations.join( df, (full_combinations["anchor"] == df["anchor"]) & (full_combinations["date_id"] == df["date_id"]), "left_outer" ).select( full_combinations["date_id"], full_combinations["anchor"], df["subs_ids"] ) # 定义窗口:按anchor分区,按date_id升序排列 fill_window = Window.partitionBy("anchor").orderBy("date_id") # 使用last函数向前填充非空值,ignoreNulls=True表示跳过空值取最近的有效数据 result = combined.withColumn( "subs_ids", F.last("subs_ids", ignoreNulls=True).over(fill_window) ).orderBy("date_id", "anchor") # 查看结果 result.show(truncate=False)
方法二:日期区间展开(高效,适合大数据量)
通过计算每个有效日期的覆盖区间直接生成完整记录,避免后续填充操作:
# 为每个anchor的有效日期计算下一个有效日期 next_date_window = Window.partitionBy("anchor").orderBy("date_id") df_with_range = df.withColumn( "next_date", F.lead("date_id", 1).over(next_date_window) ).withColumn( # 最后一条记录的next_date设为最大日期+1,确保覆盖到最后一天 "next_date", F.coalesce("next_date", F.date_add(F.max("date_id").over(Window.partitionBy()), 1)) ) # 生成当前有效日期到下一个有效日期前一天的所有日期 df_expanded = df_with_range.withColumn( "date_array", F.sequence("date_id", F.date_sub("next_date", 1), interval 1 day) ).select( "anchor", "subs_ids", F.explode("date_array").alias("date_id") ).orderBy("date_id", "anchor") # 查看结果 df_expanded.show(truncate=False)
结果说明
两种方法最终都会生成符合需求的数据集:补全所有缺失日期,每个缺失日期的subs_ids继承最近前置日期对应anchor的值,并按date_id升序排列。
内容的提问来源于stack exchange,提问作者krishna kaushik
相关产品推荐
相关产品推荐

