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

PySpark中基于前后值与前序计算结果填充空值的实现求助

用PySpark递归CTE实现依赖前后值的空值填充

当需要填充的空值依赖前一条记录的计算结果以及前后非空值时,常规Window函数(如lag/lead)无法满足需求——因为Window函数只能引用原始数据的字段,无法访问迭代计算后的新字段。此时可以用PySpark的**递归CTE(Common Table Expression)**来实现逐行递推填充。

实现步骤示例

假设你的填充规则是按行号间隔的线性插值(可根据实际规则修改计算逻辑),以下是完整实现代码:

1. 准备环境与测试数据

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, row_number, first, when
from pyspark.sql.window import Window

# 初始化SparkSession
spark = SparkSession.builder.appName("RecursiveScoreFill").getOrCreate()

# 模拟原始DataFrame
data = [
    (1, "2023-01-01", 10.0),
    (1, "2023-01-02", None),
    (1, "2023-01-03", None),
    (1, "2023-01-04", 40.0),
    (2, "2023-01-01", 20.0),
    (2, "2023-01-02", None),
    (2, "2023-01-03", 50.0)
]
df = spark.createDataFrame(data, ["id", "date", "score"])
df = df.withColumn("date", col("date").cast("date"))

2. 预处理:添加行号与后续非空值标记

为了递归时能定位前后非空值,先给每个id分组内的行按date排序并添加行号,同时标记每行之后最近的非空score及其行号:

# 定义分组排序窗口
w_group = Window.partitionBy("id").orderBy("date")
# 添加行号
df = df.withColumn("rn", row_number().over(w_group))

# 定义窗口:取当前行之后的第一个非空值
w_next_non_null = Window.partitionBy("id").orderBy("rn").rowsBetween(1, Window.unboundedFollowing)
df = df.withColumn(
    "next_score",
    when(col("score").isNull(), first(col("score"), ignorenulls=True).over(w_next_non_null)).otherwise(col("score"))
).withColumn(
    "next_rn",
    when(col("score").isNull(), first(col("rn"), ignorenulls=True).over(w_next_non_null)).otherwise(col("rn"))
)

3. 递归CTE实现填充

递归CTE分为两部分:

  • 基础CTE:取每个id分组的第一行,初始化filled_score
  • 递归部分:逐行连接前一行的计算结果,根据规则填充当前行的空值
# 注册临时表供SQL使用
df.createOrReplaceTempView("score_table")

# 执行递归CTE
filled_df = spark.sql("""
WITH RECURSIVE fill_cte AS (
    -- 基础部分:每个分组的第一行
    SELECT id, date, score, rn, next_score, next_rn, score AS filled_score
    FROM score_table
    WHERE rn = 1
    UNION ALL
    -- 递归部分:逐行计算
    SELECT
        curr.id, curr.date, curr.score, curr.rn, curr.next_score, curr.next_rn,
        CASE
            -- 非空值直接保留
            WHEN curr.score IS NOT NULL THEN curr.score
            -- 空值按线性插值计算(可替换为你的自定义规则)
            ELSE prev.filled_score + (curr.next_score - prev.filled_score) / (curr.next_rn - prev.rn) * (curr.rn - prev.rn)
        END AS filled_score
    FROM fill_cte prev
    JOIN score_table curr ON prev.id = curr.id AND curr.rn = prev.rn + 1
)
-- 输出最终结果
SELECT id, date, score, filled_score
FROM fill_cte
ORDER BY id, rn
""")

filled_df.show()

关键说明

  • 递归逻辑的核心是每一行的filled_score依赖前一行的计算结果,这是Window函数无法做到的
  • 你可以修改CASE语句中的计算逻辑,适配你的具体填充规则(比如基于日期差的插值、固定步长填充等)
  • 递归CTE要求Spark版本≥2.1,若数据量极大且递归深度过高,需注意性能优化(比如拆分大分组)

内容的提问来源于stack exchange,提问作者FerCTRO

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 22:20:38