【问题标题】:Eliminate duplicates (deduplication) in Streaming DataFrame消除 Streaming DataFrame 中的重复项(去重)
【发布时间】:2021-11-22 05:47:18
【问题描述】:

我有一个 Spark 流处理器。 数据框 dfNewExceptions 有重复项(由“ExceptionId”重复)。 由于这是一个流式数据集,因此以下查询失败:

val dfNewUniqueExceptions = dfNewExceptions.sort(desc("LastUpdateTime"))
                                    .coalesce(1)
                                    .dropDuplicates("ExceptionId")
                                    
val dfNewExceptionCore = dfNewUniqueExceptions.select("ExceptionId", "LastUpdateTime")
    dfNewExceptionCore.writeStream
      .format("console")
//      .outputMode("complete")
      .option("truncate", "false")
      .option("numRows",5000)
      .start()
      .awaitTermination(1000)
  

** 线程“主”org.apache.spark.sql.AnalysisException 中的异常:流数据帧/数据集不支持排序,除非它在完整输出模式下聚合的数据帧/数据集上;; **

这也记录在这里:https://home.apache.org/~pwendell/spark-nightly/spark-branch-2.0-docs/latest/structured-streaming-programming-guide.html

关于如何从 dfNewExceptions 中删除重复项有什么建议吗?

【问题讨论】:

    标签: scala apache-spark apache-spark-sql spark-structured-streaming delta-lake


    【解决方案1】:

    我建议遵循Streaming Deduplication 的结构化流式传输指南中解释的方法。上面写着:

    您可以使用事件中的唯一标识符对数据流中的记录进行重复数据删除。这与使用唯一标识符列的静态重复数据删除完全相同。该查询将存储来自先前记录的必要数量的数据,以便它可以过滤重复记录。与聚合类似,您可以使用带或不带水印的重复数据删除。

    带水印 - 如果重复记录的到达时间有上限,那么您可以在事件时间列上定义水印并使用 guid 和事件时间列进行重复数据删除.查询将使用水印从过去的记录中删除旧的状态数据,这些数据预计不会再有任何重复。这限制了查询必须保持的状态量。

    还给出了一个Scala中的例子:

    val dfExceptions = spark.readStream. ... // columns: ExceptionId, LastUpdateTime, ... 
    
    dfExceptions 
      .withWatermark("LastUpdateTime", "10 seconds") 
      .dropDuplicates("ExceptionId", "LastUpdateTime")
    

    【讨论】:

    • 谢谢。 dropDuplicates 用于消除重复项。但我想消除重复 - 通过采用最新的“LastUpdateTime”来删除旧事件。看起来不可能在流式传输中实现它(即使在同一个块中)?
    • 问题是,LastUpdateTime 的顺序是随机的。因此,如果不对块进行排序,我看不出这是如何实现的。您能否详细说明串联是什么意思。
    • 例如:从以下 2 个条目(它们在同一个块中)中,我想消除较旧的条目:0_7ED_25D59038A21F66638E570A8C53E4D8 3/2/2018 6:08:11 PM 0_7ED_25D59038A21F66638E5/8E/8C33 2018 年 7:28:21 PM 我试过这个:val dfNewUniqueExceptions = dfNewExceptions.withWatermark("LastUpdateTime", "60 seconds").dropDuplicates("ExceptionId", "LastUpdateTime") 期望是,事件与“2018 年 3 月 2 日 6 :08:11 PM”应该被删除。但我看到结果中仍然存在这两行。
    • 我想我现在明白了。您不想消除重复项,而是希望在 ExceptionId 上选择最长时间。这不是重复数据删除任务。
    • 我想提取所有具有唯一 ExceptionIds 的事件。如果存在具有重复 ExceptionIds 的事件,那么我想选择具有最新 LastUpdateTime 的唯一事件。
    【解决方案2】:

    您可以使用watermarking 在特定时间范围内删除重复项。

    【讨论】:

      猜你喜欢
      • 2023-03-10
      • 1970-01-01
      • 2016-12-23
      • 2018-04-08
      • 2020-02-27
      • 2017-09-21
      • 2018-08-07
      • 2017-08-31
      • 1970-01-01
      相关资源
      最近更新 更多