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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 19:39:04