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

如何在Spark Autoloader流数据中去重并保留最新modifiedon行?

基于Spark Autoloader的流数据去重与最新记录保留方案

需求说明

我通过Spark Autoloader(cloudFiles)读取得到Streaming DataFrame后,需要完成以下两项数据处理:

  • 删除完全重复的行:例如id为1的两条完全一致记录仅保留1条;
  • 保留同一id的最新记录:对同一id的多条记录,仅保留modifiedon字段值最大的那一行。

数据示例

输入数据

idmodifiedon
103/08/2023
103/08/2023
202/08/2023
203/08/2023

期望输出数据

idmodifiedon
103/08/2023
203/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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 18:05:57