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
相关产品推荐
相关产品推荐

