PySpark日期关联:保留全量数据并实现交易日期前向填充
解决数据集关联与前向填充问题
问题核心
要保留DF1的所有用户、所有日期,关联DF2的交易日期(无交易时为null),再按用户对交易日期做前向填充——之前的join写法逻辑错误,导致丢失无交易用户和不匹配日期,用左外连接+窗口函数就能解决。
步骤1:正确关联数据集
用左外连接(Left Outer Join)以DF1为主表,只匹配用户ID和日期一致的行,这样能完整保留DF1的所有数据,无交易记录的行自动填充null。
PySpark代码示例
# 先统一User_ID字段类型,避免关联失败 df1 = df1.withColumn("User_ID", df1["User_ID"].cast("string")) df2 = df2.withColumn("User_ID", df2["User_ID"].cast("string")) # 左外连接,保留DF1全量数据 joined_df = df1.join(df2, on=(df1.User_ID == df2.User_ID) & (df1.Date == df2.Transaction_Date), how="left") # 清理重复的User_ID字段 joined_df = joined_df.drop(df2.User_ID)
步骤2:按用户实现前向填充
用窗口函数按用户分组、日期排序,取最近的非null交易日期填充到后续的null行。
PySpark代码示例
from pyspark.sql import Window from pyspark.sql.functions import last, col # 定义窗口:按用户分组,日期升序,范围覆盖从分组第一行到当前行 window_spec = Window.partitionBy("User_ID").orderBy("Date").rowsBetween(Window.unboundedPreceding, Window.currentRow) # 执行前向填充 filled_df = joined_df.withColumn("Transaction_Date_Filled", last(col("Transaction_Date"), ignorenulls=True).over(window_spec))
关键说明
- 左外连接是基础:保证DF1的所有用户、日期都不丢失,没有交易的行Transaction_Date自动为null。
- 窗口函数实现填充:按用户分组后,按时间顺序把最近的有效交易日期“延续”到后续无交易的日期行。
内容的提问来源于stack exchange,提问作者BenjaminFranklinGates
相关产品推荐
相关产品推荐

