【问题标题】:How to select distinct rows from a Spark Window partition如何从 Spark Window 分区中选择不同的行
【发布时间】:2020-10-01 18:08:51
【问题描述】:

我有一个示例 DF,其中包含这样的重复行:

+-------------------+--------------------+----+-----------+-------+----------+
|ID                  |CL_ID               |NBR |DT         |TYP    |KEY       |
+--------------------+--------------------+----+-----------+-------+----------+
|1000031075_20190422 |10017157594301072477|10  |2019-04-24 |N      |0000000000|
|1000031075_20190422 |10017157594301072477|10  |2019-04-24 |N      |0000000000|
|1006473016_20190421 |10577157412800147475|11  |2019-04-21 |N      |0000000000|
|1006473016_20190421 |10577157412800147475|11  |2019-04-21 |N      |0000000000|
+--------------------+--------------------+----+-----------+-------+----------+

val w = Window.partitionBy($"ENCOUNTER_ID")

使用上面的 Spark Window 分区,是否可以选择不同的行?我期望输出 DF 为:

+-------------------+--------------------+----+-----------+-------+----------+
|ID                  |CL_ID               |NBR |DT         |TYP    |KEY       |
+--------------------+--------------------+----+-----------+-------+----------+
|1000031075_20190422 |10017157594301072477|10  |2019-04-24 |N      |0000000000|
|1006473016_20190421 |10577157412800147475|11  |2019-04-21 |N      |0000000000|
+--------------------+--------------------+----+-----------+-------+----------+

我不想使用DF.DISTINCTDF.DROPDUPLICATES,因为它会涉及洗牌。 我不喜欢使用滞后或超前,因为实时无法保证行的顺序。

【问题讨论】:

标签: apache-spark apache-spark-sql


【解决方案1】:

Window 函数也可以随机播放数据。因此,如果您的所有列都是重复的,那么df.dropDuplicates 将是更好的选择。如果您的用例想使用Window 函数,那么您可以使用以下方法。

scala> df.show()
+-------------------+--------------------+---+----------+---+----------+
|                 ID|               CL_ID|NBR|        DT|TYP|       KEY|
+-------------------+--------------------+---+----------+---+----------+
|1000031075_20190422|10017157594301072477| 10|2019-04-24|  N|0000000000|
|1000031075_20190422|10017157594301072477| 10|2019-04-24|  N|0000000000|
|1006473016_20190421|10577157412800147475| 11|2019-04-21|  N|0000000000|
|1006473016_20190421|10577157412800147475| 11|2019-04-21|  N|0000000000|
+-------------------+--------------------+---+----------+---+----------+

//You can use column in partitionBy  that need to check for duplicate and also use respective orderBy also as of now I have use sample Window

scala> val W  = Window.partitionBy(col("ID"),col("CL_ID"),col("NBR"),col("DT"), col("TYP"), col("KEY")).orderBy(lit(1))

scala> df.withColumn("duplicate", when(row_number.over(W) === lit(1), lit("Y")).otherwise(lit("N")))
              .filter(col("duplicate") === lit("Y"))
              .drop("duplicate")
              .show()
+-------------------+--------------------+---+----------+---+----------+
|                 ID|               CL_ID|NBR|        DT|TYP|       KEY|
+-------------------+--------------------+---+----------+---+----------+
|1000031075_20190422|10017157594301072477| 10|2019-04-24|  N|0000000000|
|1006473016_20190421|10577157412800147475| 11|2019-04-21|  N|0000000000|
+-------------------+--------------------+---+----------+---+----------+

【讨论】:

    【解决方案2】:

    使用大数据很好地扩展您的问题的答案:

    df.dropDuplicates(include your key cols here = ID in this case). 
    

    窗口函数对数据进行混洗,但如果您有重复的条目并想选择保留哪一个,或者想要对重复项的值求和,那么窗口函数就是您的选择

    w = Window.PartitionBy('id')
    df.agg(first( value col ).over(w)) #you can use max, min, sum, first, last depending on how you want to treat duplicates
    

    如果您想保留重复值的第三种有趣的可能性(用于记录) 之前是下面

    df.withColumn('dup_values',  collect(value_col).over(w)) 
    

    这将创建一个额外的列,每行包含一个数组,以在您删除行后保留重复值

    【讨论】:

      猜你喜欢
      • 2012-12-16
      • 2018-09-17
      • 2019-07-27
      • 2020-02-11
      • 1970-01-01
      • 2014-11-25
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多