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

Spark中按多列分组并取每组最新日期数据的最优实现方案

在Spark中实现分组取最新记录的最优方案

嘿,这个需求在Spark处理业务数据时非常常见,最优的实现方式是使用窗口函数(Window Functions),它相比传统的分组聚合再关联的方式,能减少不必要的Shuffle操作,在大数据量场景下性能更优异。下面我会一步步拆解实现过程,并给出代码示例。

核心思路

我们需要:

  • 按life id、policy id、benefit id三个字段进行分区(对应Group By的逻辑)
  • 在每个分区内,按date of commencment字段降序排序(这样最新的日期会排在最前面)
  • 为每条记录添加一个排序后的行号,然后过滤出每个分区内行号为1的记录,也就是该组的最新数据

代码实现(Scala版本)

首先加载CSV数据并处理日期格式(注意示例数据里有01/02/203这种不规范的日期,实际使用时要确保日期格式统一,或者在解析时指定兼容的格式):

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

val spark = SparkSession.builder().appName("LatestRecordByGroup").getOrCreate()

// 加载CSV数据,指定表头
val df = spark.read
  .option("header", "true")
  .option("inferSchema", "false") // 手动指定Schema更可靠,这里先自动推断
  .csv("path/to/your/data.csv")

// 转换日期字段为Date类型(根据实际日期格式调整,比如dd/MM/yyyy)
val dfWithDate = df.withColumn("date_of_commencment", to_date(col("date of commencment"), "dd/MM/yyyy"))

// 定义窗口规范:按三个字段分区,按日期降序排序
val windowSpec = Window.partitionBy("life id", "policy id", "benefit id").orderBy(col("date_of_commencment").desc)

// 添加行号并过滤出每组的第一条(最新)记录
val latestDf = dfWithDate
  .withColumn("row_num", row_number().over(windowSpec))
  .filter(col("row_num") === 1)
  .drop("row_num") // 移除临时行号字段

// 查看结果
latestDf.show()

代码实现(Python版本)

如果你用PySpark,逻辑是完全一致的:

from pyspark.sql import SparkSession
from pyspark.sql.window import Window
from pyspark.sql.functions import col, to_date, row_number

spark = SparkSession.builder.appName("LatestRecordByGroup").getOrCreate()

# 加载CSV数据
df = spark.read \
    .option("header", "true") \
    .option("inferSchema", "false") \
    .csv("path/to/your/data.csv")

# 转换日期字段
df_with_date = df.withColumn("date_of_commencment", to_date(col("date of commencment"), "dd/MM/yyyy"))

# 定义窗口规范
window_spec = Window.partitionBy("life id", "policy id", "benefit id").orderBy(col("date_of_commencment").desc())

# 获取每组最新记录
latest_df = df_with_date \
    .withColumn("row_num", row_number().over(window_spec)) \
    .filter(col("row_num") == 1) \
    .drop("row_num")

# 展示结果
latest_df.show()

为什么窗口函数是最优选择?

  • 性能更优:窗口函数可以在一次Shuffle操作中完成分区和排序,而如果用groupBy+max(date)再关联原表的方式,会产生两次Shuffle(一次分组聚合,一次关联),在大数据量下性能差距明显。
  • 逻辑更清晰:直接通过窗口函数的排序和行号过滤,就能精准获取每组的最新记录,代码可读性更强,也更容易维护。
  • 灵活性高:如果后续需要获取每组的前N条记录,只需要修改过滤条件即可,无需大幅调整代码结构。

注意事项

  • 确保日期字段的格式正确,to_date函数的格式参数要和你的数据格式匹配(比如示例中的dd/MM/yyyy),如果有不规范的日期(比如01/02/203),可以用try_cast来避免解析失败导致的报错。
  • 如果你的分组字段存在空值,窗口函数会把空值作为一个单独的分区处理,这符合SQL的分组逻辑,如果需要特殊处理空值,要提前对字段进行清洗。

内容的提问来源于stack exchange,提问作者Wajdi Ben Abderrahim

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:32:16