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

如何按版本取最新值并去空值聚合Spark DataFrame?

解决Spark DataFrame按Key聚合取各列最新版本非空值的问题

需求概述

针对指定key列,将稀疏的Spark DataFrame聚合为单条记录:

  • 移除version列
  • 各列保留对应key下最新版本(最高version值)的非空值,若高版本列值为空,则 fallback 到低版本的非空值(如示例中B列仅版本1有非空值,保留B1)

输入示例

keyversionABC
Key11A1NullNull
Key11NullB1Null
Key11NullNullC1
key12A2NullNull
key12NullNullC2

预期输出

keyABC
Key1A2B1C2

解决方案实现

以下提供两种简洁高效的实现方式:

方式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()

关键逻辑说明

  1. key格式统一:通过upper函数统一key的大小写,避免因大小写差异导致分组错误
  2. 版本优先级:按version降序排序,确保高版本数据优先被筛选
  3. 非空值 fallback:通过ignoreNulls = true或自定义UDF的空值判断,实现高版本为空时自动取低版本非空值的逻辑

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 16:32:26