如何在Spark Scala中按条件合并DataFrame的多行数据
在Spark Scala中实现按Rowkey保留最新时间戳并填充Null值
实现思路
- 按
Rowkey分组,将每个分组内的数据按timestamp降序排序,确保最新数据排在最前。 - 对除
Rowkey和timestamp外的所有列,使用窗口函数从当前行之后的旧数据中取第一个非Null值,填充当前行的Null。 - 筛选每个
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
相关产品推荐
相关产品推荐

