Spark Scala按分组最大时间戳替换列值的更优实现方案
最优实现方案:使用Spark窗口函数
你可以直接使用Spark的窗口函数完成分组最大值计算并原地替换字段值,无需额外的分组聚合+join操作,性能和可读性都更好。
完整实现代码
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ // 定义窗口规则:按指定两列分组 val groupWindow = Window.partitionBy("AUFTRAG", "AUFTRAG_POS") val resultDf = df1 // 先将字符串时间转为timestamp类型 .withColumn("ERST_TIMESTAMP", to_timestamp(col("ERST_TIMESTAMP"),"yyyy.MM.dd HH:mm:ss")) // 直接计算同分组的最大时间,覆盖原字段 .withColumn("ERST_TIMESTAMP", max("ERST_TIMESTAMP").over(groupWindow)) // 查看结果 resultDf.show()
方案优势
- 避免了额外的中间表生成和join操作,减少shuffle开销,数据量级越大性能收益越明显
- 代码语义和需求直接对齐,可读性更高,后续维护成本更低
如果后续你还需要去除同一分组的重复行,只保留每行的唯一记录,只需要在上述代码最后加上dropDuplicates("AUFTRAG", "AUFTRAG_POS", "IID_CODE")即可。
内容的提问来源于stack exchange,提问作者Arij SEDIRI
相关产品推荐
相关产品推荐

