【问题标题】:how to select first n row items based on multiple conditions in pyspark如何根据pyspark中的多个条件选择前n行项目
【发布时间】:2020-06-20 16:51:41
【问题描述】:

现在我有这样的数据:

+----+----+
|col1|   d|
+----+----+
|   A|   4|
|   A|  10|
|   A|   3|
|   B|   3|
|   B|   6|
|   B|   4|
|   B| 5.5|
|   B|  13|
+----+----+

col1 是StringType,d 是TimestampType,这里我用DoubleType 代替。 我想根据条件元组生成数据。 给定一个元组[(A,3.5),(A,8),(B,3.5),(B,10)] 我想要这样的结果

+----+---+
|col1|  d|
+----+---+
|   A|  4|
|   A| 10|
|   B|  4|
|   B| 13|
+----+---+

即对于元组中的每个元素,我们从 pyspark 数据帧中选择 d 大于元组数且 col1 等于元组字符串的前 1 行。 我已经写的是:

df_res=spark_empty_dataframe    
for (x,y) in tuples:
         dft=df.filter(df.col1==x).filter(df.d>y).limit(1)
         df_res=df_res.union(dft)

但是我认为这可能有效率问题,我不知道我是否正确。

【问题讨论】:

  • @anky 很抱歉数据让您感到困惑。我已经编辑了我的问题表,数据过滤只是关于 d 和 col1,与其他列无关。大声笑

标签: python dataframe pyspark


【解决方案1】:

避免循环的一种可能方法是从您作为输入的元组创建数据框:

t = [('A',3.5),('A',8),('B',3.5),('B',10)]
ref=spark.createDataFrame([(i[0],float(i[1])) for i in t],("col1_y","d_y"))

然后我们可以在条件下加入输入数据帧(df),然后对元组的键和值进行分组,这些键和值将重复以获得每个组的第一个值,然后删除额外的列:

(df.join(ref,(df.col1==ref.col1_y)&(df.d>ref.d_y),how='inner').orderBy("col1","d")

.groupBy("col1_y","d_y").agg(F.first("col1").alias("col1"),F.first("d").alias("d"))

.drop("col1_y","d_y")).show()

+----+----+
|col1|   d|
+----+----+
|   A|10.0|
|   A| 4.0|
|   B| 4.0|
|   B|13.0|
+----+----+

注意,如果数据帧的顺序很重要,您可以使用monotonically_increasing_id 分配一个索引列并将它们包含在聚合中,然后按索引列排序。

以另一种方式编辑而不是直接使用min 获得first

(df.join(ref,(df.col1==ref.col1_y)&(df.d>ref.d_y),how='inner')

.groupBy("col1_y","d_y").agg(F.min("col1").alias("col1"),F.min("d").alias("d"))

.drop("col1_y","d_y")).show()

+----+----+
|col1|   d|
+----+----+
|   B| 4.0|
|   B|13.0|
|   A| 4.0|
|   A|10.0|
+----+----+

【讨论】:

  • 谢谢,真的帮了我很多。
猜你喜欢
  • 1970-01-01
  • 2022-11-01
  • 1970-01-01
  • 2021-10-25
  • 1970-01-01
  • 2021-09-01
  • 1970-01-01
  • 1970-01-01
  • 2018-01-15
相关资源
最近更新 更多