如何在Spark Autoloader流数据中去重并保留最新modifiedon行?
基于Spark Autoloader的流数据去重与最新记录保留方案
需求说明
我通过Spark Autoloader(cloudFiles)读取得到Streaming DataFrame后,需要完成以下两项数据处理:
- 删除完全重复的行:例如id为1的两条完全一致记录仅保留1条;
- 保留同一id的最新记录:对同一id的多条记录,仅保留
modifiedon字段值最大的那一行。
数据示例
输入数据
| id | modifiedon |
|---|---|
| 1 | 03/08/2023 |
| 1 | 03/08/2023 |
| 2 | 02/08/2023 |
| 2 | 03/08/2023 |
期望输出数据
| id | modifiedon |
|---|---|
| 1 | 03/08/2023 |
| 2 | 03/08/2023 |
当前流数据读取代码
df = ( spark.readStream .option("delimiter", ",") .option("quote", '"') .option("mode", "permissive") .option("lineSep", "\r\n") .option("multiLine", "true") .format("cloudFiles") .option("cloudFiles.format", source_format) .option("cloudFiles.schemaLocation", checkpoint_directory) .option("header", "false") .option("escape", '"') .schema(schema) .load(data_source) )
解决方案
步骤1:删除完全重复行
直接使用dropDuplicates()方法删除所有字段完全一致的重复行:
# 移除完全重复的记录 df_dedup = df.dropDuplicates()
步骤2:保留每个id的最新记录
使用窗口函数对每个id分组,按modifiedon倒序排名,筛选出排名第一的最新记录。同时为了避免流处理中状态无限膨胀,建议添加水印(根据业务场景调整水印时长):
from pyspark.sql.window import Window from pyspark.sql.functions import col, rank # 设置水印,基于modifiedon字段,这里设置为1天(可根据实际业务调整) df_watermarked = df_dedup.withWatermark("modifiedon", "1 day") # 定义窗口:按id分组,modifiedon倒序排列 window_spec = Window.partitionBy("id").orderBy(col("modifiedon").desc()) # 计算排名并筛选最新记录 df_latest = df_watermarked.withColumn("rank", rank().over(window_spec)) \ .filter(col("rank") == 1) \ .drop("rank")
完整处理代码整合
from pyspark.sql.window import Window from pyspark.sql.functions import col, rank # 读取流数据(原有代码) df = ( spark.readStream .option("delimiter", ",") .option("quote", '"') .option("mode", "permissive") .option("lineSep", "\r\n") .option("multiLine", "true") .format("cloudFiles") .option("cloudFiles.format", source_format) .option("cloudFiles.schemaLocation", checkpoint_directory) .option("header", "false") .option("escape", '"') .schema(schema) .load(data_source) ) # 1. 移除完全重复行 df_dedup = df.dropDuplicates() # 2. 设置水印(可选但推荐,控制状态大小) df_watermarked = df_dedup.withWatermark("modifiedon", "1 day") # 3. 筛选每个id的最新记录 window_spec = Window.partitionBy("id").orderBy(col("modifiedon").desc()) df_latest = df_watermarked.withColumn("rank", rank().over(window_spec)) \ .filter(col("rank") == 1) \ .drop("rank") # 后续可将df_latest用于流输出(如写入Delta Lake、Kafka等) # df_latest.writeStream...
说明
- 水印设置:如果是无限流场景,必须添加水印来自动清理过期状态,防止内存溢出;水印时长需根据业务中数据延迟的实际情况设定。
- 窗口函数:
rank()会处理同一id下modifiedon相同的情况(如示例中id=1的两条记录),确保仅保留一条。
内容的提问来源于stack exchange,提问作者Greencolor
相关产品推荐
相关产品推荐

