【问题标题】:Removing duplicate rows based on Java DataFrame [duplicate]基于Java DataFrame删除重复行[重复]
【发布时间】:2018-07-15 12:08:28
【问题描述】:

我有一个包含以下详细信息的 DataFrame。

|id|Name|Country|version|
|1 |Jack|UK     |new    |
|1 |Jack|USA    |old    |
|2 |Rose|Germany|new    |
|3 |Sam |France |old    |

我想创建一个 DataFrame,如果数据基于“id”重复,它会选择 new version 而不是 old 版本如此

|id|Name|Country|version|
|1 |Jack|UK     |new    |
|2 |Rose|Germany|new    |
|3 |Sam |France |old    |

在 Java/Spark 中执行此操作的最佳方法是什么,或者我是否必须使用某种嵌套 SQL 查询?

简化的 SQL 版本如下所示:

WITH new_version AS (
    SELECT
      ad.id
      ,ad.name
      ,ad.country
      ,ad.version
    FROM allData ad
    WHERE ad.version = 'new'
),
old_version AS (
    SELECT
      ad.id
      ,ad.name
      ,ad.country
      ,ad.version
    FROM allData ad
    LEF JOIN new_version nv on nv.id = ad.id
    WHERE ad.version = 'old'
      AND nv.id is null
),

SELECT id, name, country, version FROM new_version
UNION ALL
SELECT id, name, country, version FROM old_version

【问题讨论】:

  • 所以你最多只有一个“新”记录和一个“旧”记录?不可能有几个“老”?

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


【解决方案1】:

假设你有一个dataframe

+---+----+-------+-------+
|id |Name|Country|version|
+---+----+-------+-------+
|1  |Jack|UK     |new    |
|1  |Jack|USA    |old    |
|2  |Rose|Germany|new    |
|3  |Sam |France |old    |
+---+----+-------+-------+

使用

创建
val df = Seq(
  ("1","Jack","UK","new"),
  ("1","Jack","USA","old"),
  ("2","Rose","Germany","new"),
  ("3","Sam","France","old")
).toDF("id","Name","Country","version")

您可以实现您对给定 sql 查询 的要求,因为 删除所有重复的 id 行,old 作为 version 使用Windowrankfilterdrop 函数如下

import org.apache.spark.sql.expressions._
def windowSpec = Window.partitionBy("id").orderBy("version")
import org.apache.spark.sql.functions._
df.withColumn("rank", rank().over(windowSpec))
  .filter(!(col("version") === "old" && col("rank") > 1))
  .drop("rank")

你应该得到最终的dataframe

+---+----+-------+-------+
|id |Name|Country|version|
+---+----+-------+-------+
|3  |Sam |France |old    |
|1  |Jack|UK     |new    |
|2  |Rose|Germany|new    |
+---+----+-------+-------+

【讨论】:

  • 似乎我无法在 spark 1.6 中使用窗口函数来绕过排名函数?
  • 对于这种情况,您应该使用 groupBy、orderBy 和 agg
【解决方案2】:

对于旧版本的 Spark,您可以将 orderBygroupBy 结合使用。根据对此question 的回答,如果数据帧在该列之后排序,则该顺序应保留在groupBy 之后。因此,以下应该有效(注意orderByidversion 列上):

val df2 = df.orderBy("id", "version")
  .groupBy("id")
  .agg(first("Name").as("Name"), first("Country").as("Country"), first("version").as("version"))

这将产生以下结果

+---+----+-------+-------+
| id|Name|Country|version|
+---+----+-------+-------+
|  3| Sam| France|    old|
|  1|Jack|     UK|    new|
|  2|Rose|Germany|    new|
+---+----+-------+-------+

【讨论】:

    猜你喜欢
    • 2016-05-31
    • 2022-01-24
    • 1970-01-01
    • 1970-01-01
    • 2019-02-26
    • 2017-04-05
    • 2016-02-21
    • 2020-01-27
    相关资源
    最近更新 更多