Spark DataFrame前向填充:同ID下state填充及action空值替换
解决方案
针对大规模Spark DataFrame的全外连接后数据填充需求,以下是高效实现方案:
完整代码实现
from pyspark.sql import Window from pyspark.sql.functions import last, coalesce, array import pandas as pd # 初始化原始DataFrame df1 = spark.createDataFrame(pd.DataFrame([[1,5,[1,2]],[1,15,[1,3]],[2,4,[3,4]]], columns=["id","day","state"])) df2 = spark.createDataFrame(pd.DataFrame([[1,10,[5,6]],[1,12,[7]],[2,4,[3]], [2,6,[10,12]],[2,10,[8,9]]], columns=["id","day","action"])) # 执行全外连接并按id、day排序 joined_df = df1.join(df2, on=["id","day"], how="fullouter").orderBy("id","day") # 定义窗口:按id分区,按day顺序处理 window_spec = Window.partitionBy("id").orderBy("day").rowsBetween(Window.unboundedPreceding, Window.currentRow) # 填充state列的null值,替换action列的null值 result_df = joined_df.withColumn( "state", last("state", ignoreNulls=True).over(window_spec) ).withColumn( "action", coalesce("action", array(0)) ) # 查看结果 result_df.show()
关键步骤说明
- 全外连接与排序:先完成全外连接,再按
id和day排序,确保后续填充逻辑按时间顺序执行。 - 窗口函数填充state:
- 用
partitionBy("id")限定仅在同一id分组内处理数据 last("state", ignoreNulls=True)会忽略null值,取当前行之前最近的非null state值,实现向下填充效果- 窗口范围设置为从分区开头到当前行,保证只回溯到当前行之前的有效数据
- 用
- action列null值替换:使用
coalesce函数直接将null值替换为[0],保留原有非null数据。
该方案基于Spark分布式计算特性,无需将全量数据拉取到单节点处理,适合大规模数据集场景。
内容的提问来源于stack exchange,提问作者ironv
相关产品推荐
相关产品推荐

