Snowflake/Snowpark实现日期区间间隙填充并延续当前值
Snowflake 实现方案
核心思路是先生成所有用户+所有月末日期的完整组合,再关联变更记录,最后用窗口函数将非空的new_ind值向后填充至下一次变更。
WITH min_max_date AS ( SELECT MIN(sales_date) AS min_date, MAX(sales_date) AS max_date FROM sales ), generated_dates AS ( SELECT LAST_DAY(DATEADD('month', ROW_NUMBER() OVER (ORDER BY seq4())-1, min_date)) AS _monthend FROM TABLE(GENERATOR(ROWCOUNT => 100000)) CROSS JOIN min_max_date QUALIFY _monthend <= max_date ), all_user_dates AS ( -- 生成每个用户+每个月末日期的完整组合 SELECT g._monthend AS monthend, c.id FROM generated_dates g CROSS JOIN (SELECT DISTINCT id FROM public.customer) c ), customer_changes AS ( SELECT id, new_ind, LAST_DAY(change_date) AS change_monthend FROM public.customer ) SELECT aud.monthend, aud.id, -- 用LAST_VALUE向后填充非空值,忽略NULL LAST_VALUE(cc.new_ind IGNORE NULLS) OVER ( PARTITION BY aud.id ORDER BY aud.monthend ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS new_ind FROM all_user_dates aud LEFT JOIN customer_changes cc ON aud.id = cc.id AND aud.monthend = cc.change_monthend ORDER BY aud.id, aud.monthend;
关键说明:
all_user_datesCTE解决了原查询只返回有变更记录行的问题,确保每个用户的每个月末日期都存在一行LAST_VALUE(cc.new_ind IGNORE NULLS)会自动将最近的非空new_ind值填充到后续的NULL行中
PySpark 实现方案
逻辑和Snowflake一致,通过生成完整的日期-用户组合,再关联变更记录,最后用窗口函数填充缺失值。
from pyspark.sql import SparkSession from pyspark.sql.functions import last, col, last_day, sequence, explode, to_date, max as spark_max, min as spark_min from pyspark.sql.window import Window spark = SparkSession.builder.appName("FillMonthlyInd").getOrCreate() # 1. 获取日期范围 sales_df = spark.table("sales") date_range = sales_df.select( spark_min("sales_date").alias("min_date"), spark_max("sales_date").alias("max_date") ).collect()[0] min_date = date_range["min_date"] max_date = date_range["max_date"] # 2. 生成所有月末日期 monthly_dates = spark.sql(f""" SELECT last_day(sequence_date) AS monthend FROM ( SELECT explode(sequence(to_date('{min_date}'), to_date('{max_date}'), interval 1 month)) AS sequence_date ) """) # 3. 获取所有用户ID customer_ids = spark.table("public.customer").select("id").distinct() # 4. 生成用户+日期的完整组合 all_user_dates = monthly_dates.crossJoin(customer_ids) # 5. 处理客户变更记录,转换为月末日期 customer_changes = spark.table("public.customer").select( "id", "new_ind", last_day(col("change_date")).alias("change_monthend") ) # 6. 关联并填充缺失值 window_spec = Window.partitionBy("id").orderBy("monthend").rowsBetween(Window.unboundedPreceding, Window.currentRow) result_df = all_user_dates.join( customer_changes, (all_user_dates["id"] == customer_changes["id"]) & (all_user_dates["monthend"] == customer_changes["change_monthend"]), "left" ).select( all_user_dates["monthend"], all_user_dates["id"], last(col("new_ind"), ignorenulls=True).over(window_spec).alias("new_ind") ).orderBy("id", "monthend") result_df.show()
关键说明:
- 用
sequence+explode生成月度序列,再取last_day得到月末日期 last(col("new_ind"), ignorenulls=True)实现和Snowflake中LAST_VALUE IGNORE NULLS相同的填充逻辑- 交叉连接确保每个用户的每个月末都有对应的行
内容的提问来源于stack exchange,提问作者Dametime
相关产品推荐
相关产品推荐

