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大时才会触发更新,但遇到以下情况就会产生重复:
- 同一
smtUidNr的多条数据msgTs相同,且在不同微批中到达:因为匹配条件不满足(s.msgTs不大于t.msgTs),会被当作新记录插入; - 乱序到达的旧数据(
msgTs小于表中已有记录):同样因为匹配条件不满足,会被插入为新记录; - 微批内的
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
相关产品推荐
相关产品推荐

