You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 17:38:10