如何基于相似记录填充并扩展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 │
需求说明
- 基于相似非空记录填充各列
null值:如第3行(a, null, 123)需填充为(a,1,123);第8行(c,3,null)需填充为(c,3,567)。 - 若某行空值存在多个匹配的非空记录,需拆分该行生成多条结果:如第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
相关产品推荐
相关产品推荐

