【问题标题】:How can I choose which duplicate rows to be dropped?如何选择要删除的重复行?
【发布时间】:2016-08-04 21:00:33
【问题描述】:

我正在尝试将新数据集与旧数据集合并,我有每个表类型的主键 Seq[String],以及具有相同架构的旧数据框和新数据框。

如果主键列值匹配,我想将旧数据框中的行替换为新数据框中的行,如果它们不匹配,我想添加该行。

到目前为止我有这个:

    val finalFrame: DataFrame = oldDF.withColumn("old/new",lit("1"))
        .union(newDF.withColumn("old/new",lit("2")))
        .dropDuplicates(primaryKeySet) 

我添加了一个 1 和 2 的文字列来跟踪哪些行是哪些行,将它们合并在一起,并根据主键列名的 Seq[String] 删除重复项。这个解决方案的问题是它不允许我指定从表中删除哪些重复项,如果我可以指定删除带有“1”的重复项是最佳的,但我愿意接受替代解决方案。

【问题讨论】:

    标签: scala apache-spark dataframe apache-spark-sql


    【解决方案1】:

    在我的头上敲了一会儿,然后想出了一个窍门。我的主键是一个序列,所以不能直接带入窗口函数中的 partitionBy,所以我这样做了:

      val windowFunction = Window.partitionBy(primaryKeySet.head, primaryKeySet.tail: _*).orderBy(desc("old/new"))
      val duplicateFreeFinalDF = finalFrame.withColumn("rownum", row_number.over(windowFunction)).where("rownum = 1").drop("rownum").drop("old/new")  
    

    基本上只是使用了可变参数扩展,因此 partitionBy 会获取我的列表,然后是 rownum 窗口函数,这样我就可以确保在出现重复的情况下获取最新的副本。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-09-23
      • 2011-04-16
      • 2013-07-23
      • 1970-01-01
      • 2022-01-09
      • 1970-01-01
      • 1970-01-01
      • 2021-10-21
      相关资源
      最近更新 更多