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

如何基于相似记录填充并扩展Polars/SQL/Spark数据集空值?

基于相似记录填充空值并拆分重复匹配行的实现方案(SQL/Polars/Spark)

问题背景

现有一个含多列的数据集,以n_1、n_2、n_3三列为例,部分列存在空值与重复值,原始数据如下:

import polars as pl

data = {
    'n_1': ['a', 'b', 'a', 'c', 'd', 'b', None, 'c', 'e'],
    'n_2': [1, 2, None, 3, 1, 2, 1, 3, 5],
    'n_3': [123, 345, 123, 567, 123, 987, 123, None, 923]
}
df = pl.DataFrame(data)

原始数据表结构:

┌──────┬──────┬──────┐
│ n_1  ┆ n_2  ┆ n_3  │
│ ---  ┆ ---  ┆ ---  │
│ str  ┆ i64  ┆ i64  │
╞══════╪══════╪══════╡
│ a    ┆ 1    ┆ 123  │
│ b    ┆ 2    ┆ 345  │
│ a    ┆ null ┆ 123  │
│ c    ┆ 3    ┆ 567  │
│ d    ┆ 1    ┆ 123  │
│ b    ┆ 2    ┆ 987  │
│ null ┆ 1    ┆ 123  │
│ c    ┆ 3    ┆ null │
│ e    ┆ 5    ┆ 923  │

需求说明

  1. 基于相似非空记录填充各列null值:如第3行(a, null, 123)需填充为(a,1,123);第8行(c,3,null)需填充为(c,3,567)。
  2. 若某行空值存在多个匹配的非空记录,需拆分该行生成多条结果:如第7行(null,1,123),匹配到(a,1,123)和(d,1,123),需拆分为两行。

期望输出

┌──────┬──────┬──────┐
│ n_1  ┆ n_2  ┆ n_3  │
│ ---  ┆ ---  ┆ ---  │
│ str  ┆ i64  ┆ i64  │
╞══════╪══════╪══════╡
│ a    ┆ 1    ┆ 123  │
│ b    ┆ 2    ┆ 345  │
│ a    ┆ 1    ┆ 123  │
│ c    ┆ 3    ┆ 567  │
│ d    ┆ 1    ┆ 123  │
│ b    ┆ 2    ┆ 987  │
│ a    ┆ 1    ┆ 123  │
│ d    ┆ 1    ┆ 123  │
│ c    ┆ 3    ┆ 567  │
│ e    ┆ 5    ┆ 923  │

Polars 实现方案

核心思路:先提取所有无空值的完整参考记录,再将原始表每行与参考表匹配(非空列完全相等则匹配),最后展开匹配结果并去重。

import polars as pl

data = {
    'n_1': ['a', 'b', 'a', 'c', 'd', 'b', None, 'c', 'e'],
    'n_2': [1, 2, None, 3, 1, 2, 1, 3, 5],
    'n_3': [123, 345, 123, 567, 123, 987, 123, None, 923]
}
df = pl.DataFrame(data)

# 提取无null的完整参考记录
reference = df.drop_nulls()

# 生成匹配条件:原始行非空列需与参考行对应列相等
match_conditions = [
    pl.when(pl.col(f"left.{col}").is_not_null())
      .then(pl.col(f"left.{col}") == pl.col(f"right.{col}"))
      .otherwise(True)
    for col in df.columns
]

# 交叉连接+过滤匹配项,保留参考行的完整值并去重排序
result = (
    df.join(reference, how="cross", suffix="_right")
      .filter(pl.all(match_conditions))
      .select([pl.col(f"{col}_right").alias(col) for col in df.columns])
      .unique()
      .sort(pl.all(df.columns))
)

print(result)

原生SQL 实现方案

逻辑与Polars一致:先获取无空值的基准数据集,通过自连接匹配非空列相等的记录,最后去重排序。

假设数据表名为data_table:

WITH reference AS (
    -- 提取所有无null的完整记录
    SELECT n_1, n_2, n_3
    FROM data_table
    WHERE n_1 IS NOT NULL AND n_2 IS NOT NULL AND n_3 IS NOT NULL
)
SELECT DISTINCT r.n_1, r.n_2, r.n_3
FROM data_table t
JOIN reference r ON
    (t.n_1 IS NULL OR t.n_1 = r.n_1)
    AND (t.n_2 IS NULL OR t.n_2 = r.n_2)
    AND (t.n_3 IS NULL OR t.n_3 = r.n_3)
ORDER BY r.n_1, r.n_2, r.n_3;

Spark 实现方案

基于Spark DataFrame API,通过提取参考集、交叉连接匹配、去重排序完成需求。

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._

object FillNullBySimilar {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder().appName("FillNull").master("local[*]").getOrCreate()
    import spark.implicits._

    val data = Seq(
      ("a", 1, 123),
      ("b", 2, 345),
      ("a", null, 123),
      ("c", 3, 567),
      ("d", 1, 123),
      ("b", 2, 987),
      (null, 1, 123),
      ("c", 3, null),
      ("e", 5, 923)
    ).toDF("n_1", "n_2", "n_3")

    // 提取无null的参考记录
    val reference = data.na.drop()

    // 构建动态匹配条件
    val matchConditions = data.columns.map(col => 
      when(col(col).isNotNull, col(col) === col(s"${col}_right")).otherwise(lit(true))
    ).reduce(_ && _)

    val result = data
      .crossJoin(reference.toDF(data.columns.map(c => s"${c}_right"): _*))
      .filter(matchConditions)
      .select(data.columns.map(c => col(s"${c}_right").alias(c)): _*)
      .dropDuplicates()
      .orderBy(data.columns.map(col): _*)

    result.show()
    spark.stop()
  }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 20:27:05