【问题标题】:How to retrieve unique values in each window in pyspark dataframe如何在pyspark数据框中的每个窗口中检索唯一值
【发布时间】:2019-04-03 09:13:11
【问题描述】:

我有以下火花数据框:

from pyspark.sql import SparkSession
spark = SparkSession.builder.appName('').getOrCreate()
df = spark.createDataFrame([(1, "a", "2"), (2, "b", "2"),(3, "c", "2"), (4, "d", "2"),
                (5, "b", "3"), (6, "b", "3"),(7, "c", "2")], ["nr", "column2", "quant"])

返回我:

+---+-------+------+
| nr|column2|quant |
+---+-------+------+
|  1|      a|     2|
|  2|      b|     2|
|  3|      c|     2|
|  4|      d|     2|
|  5|      b|     3|
|  6|      b|     3|
|  7|      c|     2|
+---+-------+------+

我想检索每 3 个分组行(从窗口大小为 3 的每个窗口)中 quant 列具有唯一值的行。如下图:

这里红色是窗口大小,每个窗口我只保留 quant 唯一的绿色行:

我想得到的输出如下:

+---+-------+------+
| nr|column2|values|
+---+-------+------+
|  1|      a|     2|
|  4|      d|     2|
|  5|      b|     3|
|  7|      c|     2|
+---+-------+------+

我是 Spark 的新手,所以我将不胜感激。谢谢

【问题讨论】:

  • 你如何创建你的红色组?你的分组条件是什么?
  • 前3行为一组,后3行为下一组
  • 第一次按什么排序?
  • 从顶部开始。数据集已经排序。只是前 3 行是一组,然后下 3 行是下一组。它应该滚动
  • 分布式文件系统或数据库中没有任何内容已排序。它就像一袋弹珠,你可以把它们从袋子里分出来,但不能在里面。

标签: python group-by pyspark apache-spark-sql window


【解决方案1】:

假设将 3 条记录分组是基于 'nr' 列,这种方法应该适合您。

使用udf,决定是否选择记录,lag,用于获取上一行数据。

def tag_selected(index, current_quant, prev_quant1, prev_quant2):                                                                                                    
    if index % 3 == 1:  # first record in each group is always selected                                                                                              
        return True                                                                                                                                                  
    if index % 3 == 2 and current_quant != prev_quant1: # second record will be selected if prev quant is not same as current                                        
        return True                                                                                                                                                  
    if index % 3 == 0 and current_quant != prev_quant1 and current_quant != prev_quant2: # third record will be selected if prev quant are not same as current       
        return True                                                                                                                                                  
    return False                                                                                                                                                     

tag_selected_udf = udf(tag_selected, BooleanType())                                                                                                                  

df = spark.createDataFrame([(1, "a", "2"), (2, "b", "2"),(3, "c", "2"), (4, "d", "2"),
                (5, "b", "3"), (6, "b", "3"),(7, "c", "2")], ["nr", "column2", "quant"])

window = Window.orderBy("nr")

df = df.withColumn("prev_quant1", lag(col("quant"),1, None).over(window))\
       .withColumn("prev_quant2", lag(col("quant"),2, None).over(window)) \
       .withColumn("selected", 
                   tag_selected_udf(col('nr'),col('quant'),col('prev_quant1'),col('prev_quant2')))\
       .filter(col('selected') == True).drop("prev_quant1","prev_quant2","selected")
df.show()

结果

+---+-------+-----+
| nr|column2|quant|
+---+-------+-----+
|  1|      a|    2|
|  4|      d|    2|
|  5|      b|    3|
|  7|      c|    2|
+---+-------+-----+

【讨论】:

  • @Sascha,如果它对你有用,你能接受这个作为答案吗,谢谢。
  • 首先感谢您的帮助,我对这段代码有疑问,所以每次我们创建新列时在df.withColumn("prev_quant1", lag(col("quant"),1, None).over(window))....?这意味着在这种情况下,如果数据更大,它将占用太多内存?第二个问题是关于def tag_selected。这是否意味着它像滚动窗口一样工作,但不将 3 行(1,2,3)组和接下来的 3 行(4,5,6)作为不同的组?因为当我运行它时,我得到了 1、4、5、6、7 的输出。是否可以将窗口大小定义为 3?
  • 对 rdd/df/datasets 的转换将创建新的 rdds 作为 rdd 不可变。所以添加新的col会创建新的rdd实例,但是spark内存管理决定旧的rdds是否需要保留或清除。tag_selected,像滚动窗口一样采用current,prev quants,但里面的逻辑忽略了其他窗口元素。对于第 4 行,使用第 2、3 行数据调用 udf,但 udf 会忽略,因为它根据索引知道它是组的第一行并且总是需要选择。代码,如果 index % 3 == 0...... 失败并返回 False,则不应选择 6 作为 udf 中的条件。
  • 感谢您的解释。我想再澄清一件事。如果 spark 数据帧已经按 nr 排序但 nr 不是从 1 开始怎么办。例如,nr 的顺序如下:“15.3、22.8、37.1、39、56 等”。在这种情况下,最好添加新的索引列并遵循该逻辑?还是有其他关于 nr 的选择?
  • 是的,如果 nr 值不是索引号,那么您需要添加新的 col 作为索引号的索引。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2010-10-14
  • 2021-07-07
  • 1970-01-01
  • 2021-11-29
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多