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

基于指定列合并Spark DataFrame多行 实现CDC数据变更处理

Spark同id多行合并取最新非空值实现方案

你遇到的是CDC变更数据合并的典型场景,无需多表关联,直接用Spark内置的DataFrame函数即可实现,以下是两种常用实现方式:


方案1:窗口函数 + last忽略空值(兼容Spark 2.x及以上所有版本)

核心逻辑是按id分区、按更新时间排序后,对每个业务字段取窗口内最后一个非空值,最后去重即可:

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._

// 定义窗口:按id分区,按更新时间升序排序,窗口范围覆盖同id所有行
val idWindow = Window.partitionBy("id")
  .orderBy("update_time")
  .rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing)

val resultDF = originalDF
  .select(
    col("id"),
    // last的第二个参数设为true表示忽略NULL值,取最新的非空记录
    last(col("name"), ignoreNulls = true).over(idWindow).as("name"),
    last(col("age"), ignoreNulls = true).over(idWindow).as("age"),
    last(col("city"), ignoreNulls = true).over(idWindow).as("city"),
    max(col("update_time")).over(idWindow).as("update_time")
  )
  // 同id的所有行计算后结果完全一致,直接按id去重即可
  .dropDuplicates("id")

方案2:groupBy + max_by函数(Spark 3.0+支持,性能更优)

Spark 3.0新增的max_by函数可以直接按排序字段取对应目标列的最大值,写法更简洁,shuffle开销更小:

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

val resultDF = originalDF
  .groupBy("id")
  .agg(
    // 按update_time升序,取最大时间对应的非空字段值
    max_by(col("name"), col("update_time")).as("name"),
    max_by(col("age"), col("update_time")).as("age"),
    max_by(col("city"), col("update_time")).as("city"),
    max(col("update_time")).as("update_time")
  )

两种方案执行后得到的结果和你要求的输出完全一致。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 21:45:06