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

如何在Spark Scala中按条件合并DataFrame的多行数据

在Spark Scala中实现按Rowkey保留最新时间戳并填充Null值

实现思路

  1. 按Rowkey分组,将每个分组内的数据按timestamp降序排序,确保最新数据排在最前。
  2. 对除Rowkey和timestamp外的所有列,使用窗口函数从当前行之后的旧数据中取第一个非Null值,填充当前行的Null。
  3. 筛选每个Rowkey下timestamp最大的行,得到最终结果。

完整代码示例

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

// 假设你的原始DataFrame名为df,这里替换为你的数据源读取逻辑
val df = spark.read... 

// 定义填充窗口:按Rowkey分组,timestamp降序,窗口范围为当前行之后的所有旧数据
val fillWindow = Window.partitionBy("Rowkey")
  .orderBy(desc("timestamp"))
  .rowsBetween(1, Window.unboundedFollowing)

// 获取需要填充的列(排除Rowkey和timestamp)
val fillColumns = df.columns.filter(col => !List("Rowkey", "timestamp").contains(col))

// 逐列填充Null:优先保留当前行非Null值,为Null则取后续旧行的第一个非Null值
val filledDataFrame = fillColumns.foldLeft(df) { (tempDf, colName) =>
  tempDf.withColumn(
    colName,
    coalesce(col(colName), first(col(colName), ignoreNulls = true).over(fillWindow))
  )
}

// 筛选每个Rowkey下的最新时间戳行
val maxTimestampWindow = Window.partitionBy("Rowkey")
val resultDataFrame = filledDataFrame
  .withColumn("max_ts", max(col("timestamp")).over(maxTimestampWindow))
  .where(col("timestamp") === col("max_ts"))
  .drop("max_ts")

// 查看最终结果
resultDataFrame.show()

代码说明

  • fillWindow:限定了分组和排序规则,窗口范围仅包含当前行之后的旧数据,确保填充值来自时间更早的有效数据。
  • coalesce + first:coalesce会优先保留当前行的非Null值,只有当当前行为Null时,才会从窗口内取第一个非Null值填充,完全匹配需求。
  • maxTimestampWindow:用于标记每个Rowkey的最新时间戳,最终筛选出仅包含最新数据的行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 01:25:37