PySpark下Spark DataFrame空值前后向填充问题及修正需求
Spark DataFrame 自定义填充需求实现
我需要在Spark DataFrame中创建窗口分区列desired_output,规则如下:
- 按
user_id分组,严格遵循分组内记录的原始顺序(不依赖date列的时间或值顺序) - 分组开头的空值,用后续首个非空值后向填充
- 其余位置的空值,用前面最近的非空值前向填充
初始示例Spark DataFrame
columns = ['user_id', 'date', 'desired_outcome'] data = [\ ('1', None, '2022-01-05'),\ ('1', None, '2022-01-05'),\ ('1', '2022-01-05', '2022-01-05'),\ ('1', None, '2022-01-05'),\ ('1', None, '2022-01-05'),\ ('2', None, '2022-01-07'),\ ('2', None, '2022-01-07'),\ ('2','2022-01-07', '2022-01-07'),\ ('2',None, '2022-01-07'),\ ('2','2022-01-09', '2022-01-09'),\ ('2',None, '2022-01-09'),\ ('3','2022-01-01', '2022-01-01'),\ ('3',None, '2022-01-01'),\ ('3',None, '2022-01-01'),\ ('3','2022-01-04', '2022-01-04'),\ ('3',None, '2022-01-04'),\ ('3',None, '2022-01-04')] sample_df = spark.createDataFrame(data, columns)
【更新】
基于@user238607的提案解决方案测试时发现错误:当date列值不按时间顺序排列时,计算结果不符合预期(错误值在correct列中用---标记)。补充说明:期望输出仅依赖user_id分组内的原始记录顺序,date列可为任意类型(数值、字符串等)。
测试代码
from pyspark.sql import Window from pyspark import SQLContext from pyspark.sql.functions import * import pyspark.sql.functions as F sc = spark #SparkContext('local') sqlContext = SQLContext(sc) data1 = [ ('1', None, '2022-02-12'), ('1', None, '2022-02-12'), ('1', '2022-02-12', '2022-02-12'), ('1', None, '2022-02-12'), ('1', None, '2022-02-12'), ('2', None, '2022-04-09'), ('2', None, '2022-04-09'), ('2','2022-04-09', '2022-04-09'), ('2',None, '2022-04-09'), ('2','2022-01-07', '2022-01-07'), ('2',None, '2022-01-07'), ('3','2022-11-05', '2022-11-05'), ('3',None, '2022-11-05'), ('3',None, '2022-11-05'), ('3','2022-01-04', '2022-01-04'), ('3',None, '2022-01-04'), ('3',None, '2022-01-04'), ('3','2022-04-15', '2022-04-15'), ('3',None, '2022-04-15'), ] columns = ['row_id', 'user_id', 'date', 'desired_outcome_given'] ## 新增row_id以保留原始行顺序 data2 = [ (index, item[0], item[1], item[2]) for index, item in enumerate(data1) ] df1 = sqlContext.createDataFrame(data=data2, schema=columns) print("Given dataframe") df1.show(n=10, truncate=False) window_min_spec = Window.partitionBy("user_id").orderBy(F.col("row_id").asc()).rowsBetween(0, Window.unboundedFollowing) window_max_spec = Window.partitionBy("user_id").orderBy(F.col("row_id").asc()).rowsBetween(Window.unboundedPreceding, 0) df1 = df1.withColumn("first_date", F.first("date", ignorenulls=True).over(window_min_spec)) df1 = df1.withColumn("last_date", F.last("date", ignorenulls=True).over(window_max_spec)) print("Calculated first and last dates") df1.show(truncate=False) print("Final dataframe") output = df1.\ withColumn("desired_outcome_calculated", F.least(*["first_date", "last_date"]))\ .withColumn("correct", F.when(F.col("desired_outcome_given") == F.col("desired_outcome_calculated"),F.lit('true')).otherwise('---'))\ .select('row_id', 'user_id', 'date', "desired_outcome_given", "desired_outcome_calculated", "correct") output.show(truncate=False)
测试输出
Given dataframe +------+-------+----------+---------------------+ |row_id|user_id|date |desired_outcome_given| +------+-------+----------+---------------------+ |0 |1 |NULL |2022-02-12 | |1 |1 |NULL |2022-02-12 | |2 |1 |2022-02-12|2022-02-12 | |3 |1 |NULL |2022-02-12 | |4 |1 |NULL |2022-02-12 | |5 |2 |NULL |2022-04-09 | |6 |2 |NULL |2022-04-09 | |7 |2 |2022-04-09|2022-04-09 | |8 |2 |NULL |2022-04-09 | |9 |2 |2022-01-07|2022-01-07 | +------+-------+----------+---------------------+ only showing top 10 rows Calculated first and last dates +------+-------+----------+---------------------+----------+----------+ |row_id|user_id|date |desired_outcome_given|first_date|last_date | +------+-------+----------+---------------------+----------+----------+ |0 |1 |NULL |2022-02-12 |2022-02-12|NULL | |1 |1 |NULL |2022-02-12 |2022-02-12|NULL | |2 |1 |2022-02-12|2022-02-12 |2022-02-12|2022-02-12| |3 |1 |NULL |2022-02-12 |NULL |2022-02-12| |4 |1 |NULL |2022-02-12 |NULL |2022-02-12| |5 |2 |NULL |2022-04-09 |2022-04-09|NULL | |6 |2 |NULL |2022-04-09 |2022-04-09|NULL | |7 |2 |2022-04-09|2022-04-09 |2022-04-09|2022-04-09| |8 |2 |NULL |2022-04-09 |2022-01-07|2022-04-09| |9 |2 |2022-01-07|2022-01-07 |2022-01-07|2022-01-07| |10 |2 |NULL |2022-01-07 |NULL |2022-01-07| |11 |3 |2022-11-05|2022-11-05 |2022-11-05|2022-11-05| |12 |3 |NULL |2022-11-05 |2022-01-04|2022-11-05| |13 |3 |NULL |2022-11-05 |2022-01-04|2022-11-05| |14 |3 |2022-01-04|2022-01-04 |2022-01-04|2022-01-04| |15 |3 |NULL |2022-01-04 |2022-04-15|2022-01-04| |16 |3 |NULL |2022-01-04 |2022-04-15|2022-01-04| |17 |3 |2022-04-15|2022-04-15 |2022-04-15|2022-04-15| |18 |3 |NULL |2022-04-15 |NULL |2022-04-15| +------+-------+----------+---------------------+----------+----------+ Final dataframe +------+-------+----------+---------------------+--------------------------+-------+ |row_id|user_id|date |desired_outcome_given|desired_outcome_calculated|correct| +------+-------+----------+---------------------+--------------------------+-------+ |0 |1 |NULL |2022-02-12 |2022-02-12 |true | |1 |1 |NULL |2022-02-12 |2022-02-12 |true | |2 |1 |2022-02-12|2022-02-12 |2022-02-12 |true | |3 |1 |NULL |2022-02-12 |2022-02-12 |true | |4 |1 |NULL |2022-02-12 |2022-02-12 |true | |5 |2 |NULL |2022-04-09 |2022-04-09 |true | |6 |2 |NULL |2022-04-09 |2022-04-09 |true | |7 |2 |2022-04-09|2022-04-09 |2022-04-09 |true | |8 |2 |NULL |2022-04-09 |2022-01-07 |--- | |9 |2 |2022-01-07|2022-01-07 |2022-01-07 |true | |10 |2 |NULL |2022-01-07 |2022-01-07 |true | |11 |3 |2022-11-05|2022-11-05 |2022-11-05 |true | |12 |3 |NULL |2022-11-05 |2022-01-04 |--- | |13 |3 |NULL |2022-11-05 |2022-01-04 |--- | |14 |3 |2022-01-04|2022-01-04 |2022-01-04 |true | |15 |3 |NULL |2022-01-04 |2022-01-04 |true | |16 |3 |NULL |2022-01-04 |2022-01-04 |true | |17 |3 |2022-04-15|2022-04-15 |2022-04-15 |true | |18 |3 |NULL |2022-04-15 |2022-04-15 |true | +------+-------+----------+---------------------+--------------------------+-------+
内容的提问来源于stack exchange,提问作者ShamilS
相关产品推荐
相关产品推荐

