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

Spark Structured Streaming处理Azure Databricks Delta流数据去重问题求助

解决Spark Structured Streaming + Delta Lake Upsert后重复记录的问题

首先直接回应你的疑问:Spark Structured Streaming完全支持聚合、窗口函数和orderBy子句,只是有一些适配流式场景的限制:

  • 聚合操作需要搭配对应的outputMode(比如update或complete,不能用append);
  • 基于事件时间的窗口函数需要配合**水印(Watermark)**来处理乱序数据并控制状态大小;
  • orderBy可以在微批内部使用,但全局排序因为流式的微批特性无法直接实现,通常只在每个微批内做局部排序。

接下来分析你当前代码导致重复记录的原因,并给出针对性的修改方案:

问题根源分析

你的Upsert逻辑的匹配条件是:

ON (s.smtUidNr = t.smtUidNr and s.msgTs>t.msgTs)

这个条件意味着只有当新数据的msgTs比表中已有记录的msgTs大时才会触发更新,但遇到以下情况就会产生重复:

  1. 同一smtUidNr的多条数据msgTs相同,且在不同微批中到达:因为匹配条件不满足(s.msgTs不大于t.msgTs),会被当作新记录插入;
  2. 乱序到达的旧数据(msgTs小于表中已有记录):同样因为匹配条件不满足,会被插入为新记录;
  3. 微批内的dropDuplicates依赖date_timestamp分组,只能过滤同一日期内的重复,跨日期的乱序数据还是会漏进Upsert环节。

解决方案:优化Upsert逻辑+微批内去重

我们需要调整Upsert的匹配逻辑,确保每个smtUidNr始终保留**最新(最大msgTs)**的记录,同时在微批内先处理掉重复数据:

步骤1:简化微批内的处理逻辑

用DataFrame API替代大量临时视图,让代码更清晰,同时确保每个微批内的smtUidNr只保留msgTs最大的记录:

import org.apache.spark.sql._
import org.apache.spark.sql.functions._

// 读取流数据
val df = spark.readStream.format("delta").load("abfss://abc@hjklinfo.dfs.core.windows.net/entrypacket/")

// 展开嵌套结构并提取所需字段
val entrypacketDF = df.select(
  explode(col("details")).alias("dcl"),
  explode(col("invdetails")).alias("inv"),
  explode(col("eventdetails")).alias("evt"),
  explode(col("smtdetails")).alias("smt"),
  col("msgHdr.msgTs"),
  col("msgHdr.msgInfSrcCd")
)

// 转换msgTs为timestamp类型,并按smtUidNr分组保留最新记录
val latestPerSmtUidDF = entrypacketDF
  .withColumn("msgTs", col("msgTs").cast("timestamp"))
  .groupBy("smtUidNr")
  .agg(
    max("dcl").alias("dcl"),
    max("inv").alias("inv"),
    max("evt").alias("evt"),
    max("smt").alias("smt"),
    max("msgTs").alias("msgTs"),
    max("msgInfSrcCd").alias("msgInfSrcCd")
  )

步骤2:修改Upsert逻辑,确保只保留最新记录

调整MERGE的匹配条件为仅匹配smtUidNr,然后在匹配时只更新新数据msgTs更大的情况,避免旧数据插入重复:

def upsertToDelta(microBatchOutputDF: DataFrame, batchId: Long): Unit = {
  microBatchOutputDF.createOrReplaceTempView("updates")
  
  microBatchOutputDF.sparkSession.sql("""
    MERGE INTO raw t
    USING updates s
    ON s.smtUidNr = t.smtUidNr
    WHEN MATCHED AND s.msgTs > t.msgTs THEN 
      UPDATE SET *
    WHEN NOT MATCHED THEN 
      INSERT *
  """)
}

步骤3:启动流查询

用update模式启动流,确保聚合后的结果能正确触发Upsert:

latestPerSmtUidDF
  .writeStream
  .format("delta")
  .foreachBatch(upsertToDelta _)
  .outputMode("update")
  .start()
  .awaitTermination()

可选:添加水印处理乱序数据

如果你的乱序数据有明确的延迟阈值(比如最多延迟1小时),可以添加水印来自动清理过期状态,避免状态无限增长:

val latestPerSmtUidDF = entrypacketDF
  .withColumn("msgTs", col("msgTs").cast("timestamp"))
  .withWatermark("msgTs", "1 hour") // 设置1小时的水印,丢弃超过1小时的旧数据
  .groupBy("smtUidNr")
  .agg(
    max("dcl").alias("dcl"),
    max("inv").alias("inv"),
    max("evt").alias("evt"),
    max("smt").alias("smt"),
    max("msgTs").alias("msgTs"),
    max("msgInfSrcCd").alias("msgInfSrcCd")
  )

验证效果

修改后,每个smtUidNr只会保留msgTs最大的记录:

  • 同一微批内的重复数据会被group by聚合处理;
  • 跨微批的新数据如果msgTs更大,会更新现有记录;
  • 跨微批的旧数据(msgTs更小)不会触发插入或更新,彻底避免重复。

内容的提问来源于stack exchange,提问作者AKSHAY SHINGOTE

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:15:50