如何在Databricks中用递增值填充DataFrame的空值行?
问题描述
我在Databricks中操作名为df的DataFrame,value列存在多个空值行,数据结构如下:
+----+------------+-------+ | id | date | value | +----+------------+-------+ | 1 | 2024-01-01 | 0 | | 1 | 2024-01-02 | null | | 1 | 2024-01-03 | null | | 1 | 2024-01-04 | null | | 1 | 2024-01-05 | 4 | | 1 | 2024-01-06 | null | | 2 | 2024-03-20 | 0 | | 2 | 2024-03-21 | null | | 2 | 2024-03-22 | null | | 2 | 2024-03-23 | 3 | | 2 | 2024-03-24 | null | | 2 | 2024-03-25 | 0 | | 2 | 2024-03-26 | null | | 2 | 2024-03-27 | null | | 2 | 2024-03-28 | null | | 2 | 2024-03-29 | 4 | +----+------------+-------+
需要按id分区、按date排序,将空值替换为前一行的值加1(第一行保留原值,无论是否为空)。
伪代码算法
for (int i = 1; i < number_of_rows; i++) { if value[i] == null { value[i] = value[i-1] + 1 } }
期望输出
+----+------------+-------+--------------+ | id | date | value | value_filled | +----+------------+-------+--------------+ | 1 | 2024-01-01 | 0 | 0 | | 1 | 2024-01-02 | null | 1 | | 1 | 2024-01-03 | null | 2 | | 1 | 2024-01-04 | null | 3 | | 1 | 2024-01-05 | 4 | 4 | | 1 | 2024-01-06 | null | 5 | | 2 | 2024-03-20 | 0 | 0 | | 2 | 2024-03-21 | null | 1 | | 2 | 2024-03-22 | null | 2 | | 2 | 2024-03-23 | 3 | 3 | | 2 | 2024-03-24 | null | 4 | | 2 | 2024-03-25 | 0 | 0 | | 2 | 2024-03-26 | null | 1 | | 2 | 2024-03-27 | null | 2 | | 2 | 2024-03-28 | null | 3 | | 2 | 2024-03-29 | 4 | 4 | +----+------------+-------+--------------+
已尝试的方法
方法1:使用last函数填充空值为前一行非空值
from pyspark.sql import functions as F from pyspark.sql.window import Window w = Window.partitionBy('id').orderBy('date').rowsBetween(Window.unboundedPreceding, 0) df = df.withColumn('value_filled', F.last('value', ignorenulls=True).over(w))
执行结果:
+----+------------+-------+--------------+ | id | date | value | value_filled | +----+------------+-------+--------------+ | 1 | 2024-01-01 | 0 | 0 | | 1 | 2024-01-02 | null | 0 | | 1 | 2024-01-03 | null | 0 | | 1 | 2024-01-04 | null | 0 | | 1 | 2024-01-05 | 4 | 4 | | 1 | 2024-01-06 | null | 4 | | 2 | 2024-03-20 | 0 | 0 | | 2 | 2024-03-21 | null | 0 | | 2 | 2024-03-22 | null | 0 | | 2 | 2024-03-23 | 3 | 3 | | 2 | 2024-03-24 | null | 3 | | 2 | 2024-03-25 | 0 | 0 | | 2 | 2024-03-26 | null | 0 | | 2 | 2024-03-27 | null | 0 | | 2 | 2024-03-28 | null | 0 | | 2 | 2024-03-29 | 4 | 4 | +----+------------+-------+--------------+
该方法只能将空值填充为最近的非空值,无法实现递增。
方法2:尝试为空值行加1
w = Window.partitionBy('id').orderBy('date').rowsBetween(Window.unboundedPreceding, 0) df = df.withColumn('value_filled', F.when(F.isnull(F.col('value')), F.last(F.col('value'), ignorenulls=True).over(w) + 1).otherwise(F.col('value')))
执行结果:
+----+------------+-------+--------------+ | id | date | value | value_filled | +----+------------+-------+--------------+ | 1 | 2024-01-01 | 0 | 0 | | 1 | 2024-01-02 | null | 1 | | 1 | 2024-01-03 | null | 1 | | 1 | 2024-01-04 | null | 1 | | 1 | 2024-01-05 | 4 | 4 | | 1 | 2024-01-06 | null | 5 | | 2 | 2024-03-20 | 0 | 0 | | 2 | 2024-03-21 | null | 1 | | 2 | 2024-03-22 | null | 1 | | 2 | 2024-03-23 | 3 | 3 | | 2 | 2024-03-24 | null | 4 | | 2 | 2024-03-25 | 0 | 0 | | 2 | 2024-03-26 | null | 1 | | 2 | 2024-03-27 | null | 1 | | 2 | 2024-03-28 | null | 1 | | 2 | 2024-03-29 | 4 | 4 | +----+------------+-------+--------------+
该方法仅能将每组空值的第一个行加1,后续空值无法持续递增。
解决方案
要实现空值按前一行值递增的效果,需要先以非空值为分界对数据分组,再在组内计算偏移量:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 步骤1:按id分区、date排序,标记每个非空value所属的组 w_group = Window.partitionBy('id').orderBy('date') df = df.withColumn('group_id', F.sum(F.when(F.col('value').isNotNull(), 1)).over(w_group)) # 步骤2:在每个id+group_id的组内,计算当前行与组内首行的偏移量 w_row = Window.partitionBy('id', 'group_id').orderBy('date') df = df.withColumn('row_offset', F.row_number().over(w_row) - 1) # 步骤3:获取每组的基准值(组内第一个非空value),加上偏移量得到填充值 w_base = Window.partitionBy('id', 'group_id') df = df.withColumn('base_value', F.first('value', ignorenulls=True).over(w_base)) df = df.withColumn('value_filled', F.col('base_value') + F.col('row_offset')) # 可选:删除中间辅助列 df = df.drop('group_id', 'row_offset', 'base_value')
执行后即可得到符合要求的输出结果。
内容的提问来源于stack exchange,提问作者Mikey Antonakakis
相关产品推荐
相关产品推荐

