如何按版本取最新值并去空值聚合Spark DataFrame?
解决Spark DataFrame按Key聚合取各列最新版本非空值的问题
需求概述
针对指定key列,将稀疏的Spark DataFrame聚合为单条记录:
- 移除
version列 - 各列保留对应
key下最新版本(最高version值)的非空值,若高版本列值为空,则 fallback 到低版本的非空值(如示例中B列仅版本1有非空值,保留B1)
输入示例
| key | version | A | B | C |
|---|---|---|---|---|
| Key1 | 1 | A1 | Null | Null |
| Key1 | 1 | Null | B1 | Null |
| Key1 | 1 | Null | Null | C1 |
| key1 | 2 | A2 | Null | Null |
| key1 | 2 | Null | Null | C2 |
预期输出
| key | A | B | C |
|---|---|---|---|
| Key1 | A2 | B1 | C2 |
解决方案实现
以下提供两种简洁高效的实现方式:
方式1:窗口函数 + 去重(推荐)
利用窗口函数按key分区、version降序排序,对每列取第一个非空值后去重,得到每个key的唯一记录:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions.{coalesce, first, upper} // 假设原始DataFrame名为df val windowSpec = Window.partitionBy(upper($"key").alias("key_normalized")) .orderBy($"version".desc) val result = df // 统一key大小写,处理示例中Key1和key1的分组问题 .withColumn("key", upper($"key")) // 按窗口取每列第一个非空值,自动跳过空值优先取高版本数据 .withColumn("A", first(coalesce($"A"), ignoreNulls = true).over(windowSpec)) .withColumn("B", first(coalesce($"B"), ignoreNulls = true).over(windowSpec)) .withColumn("C", first(coalesce($"C"), ignoreNulls = true).over(windowSpec)) // 筛选目标列并去重,保留每个key的唯一聚合记录 .select("key", "A", "B", "C") .distinct() result.show()
方式2:分组聚合 + 自定义UDF
如果需要更灵活的逻辑,可自定义UDF收集各版本的非空值,再取最新的一个:
import org.apache.spark.sql.functions.{udf, collect_list, struct} // 自定义UDF:从按version降序排列的结构列表中,提取第一个非空值 val getLatestNonNull = udf((values: Seq[(Int, String)]) => { values.sortBy(-_._1).find(_._2 != null).map(_._2).orNull }) val result = df .withColumn("key", upper($"key")) // 为每列构造(version, 列值)的结构,便于后续按版本排序筛选 .withColumn("A_struct", struct($"version", $"A")) .withColumn("B_struct", struct($"version", $"B")) .withColumn("C_struct", struct($"version", $"C")) // 按key分组,收集各列的版本-值结构列表,再用UDF提取最新非空值 .groupBy("key") .agg( getLatestNonNull(collect_list($"A_struct")).alias("A"), getLatestNonNull(collect_list($"B_struct")).alias("B"), getLatestNonNull(collect_list($"C_struct")).alias("C") ) result.show()
关键逻辑说明
- key格式统一:通过
upper函数统一key的大小写,避免因大小写差异导致分组错误 - 版本优先级:按
version降序排序,确保高版本数据优先被筛选 - 非空值 fallback:通过
ignoreNulls = true或自定义UDF的空值判断,实现高版本为空时自动取低版本非空值的逻辑
内容的提问来源于stack exchange,提问作者karas27
相关产品推荐
相关产品推荐

