基于指定列合并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
相关产品推荐
相关产品推荐

