【问题标题】:pyspark create n columns based on selecting column values from other rows in the same grouppyspark 根据从同一组中的其他行中选择列值创建 n 列
【发布时间】:2021-03-02 00:21:44
【问题描述】:

我想根据以下条件向数据框的每一行添加 n 列
1) 对给定的按列分组的操作应用分组
2)处理每个组
3)对于组中的每条记录,根据给定条件从其他记录中选择一列的最后 n 个值

数据框:

   demo,origin,gap,date_nor,date,ratings
   men_20_21, india, -1, 0, 1/11/2020,0.1
   men_20_21, india, 0, 1, 2/11/2020,0.2
   men_20_21, india, 1, 2, 3/11/2020,0.3
   men_20_21, india, 2, 3, 4/11/2020,0.4
   men_20_21, india, 3, 4, 5/11/2020,0.5
   men_30_35, india, 4, 5, 6/11/2020,0.6
   men_30_35, india, 5, 6, 7/11/2020,0.7
   men_30_35, india, 6, 7, 8/11/2020,0.8
   men_30_35, india, 7, 8, 9/11/2020,0.9
   men_30_35, india, 8, 9, 10/11/2020,0.10
   men_30_35, india, 9, 10, 11/11/2020,0.11
   men_30_35, india, 10, 11, 12/11/2020,0.12

以上数据框按 demo 和 origin 列分组

输出数据帧

   demo,origin,gap,date_nor,date,ratings,1last_rating,2last_rating,3last_rating,4last_rating,5last_rating
   men_20_21, india, -1, 0, 1/11/2020,0.1,null,null,null,null,null
   men_20_21, india, 0, 1, 2/11/2020,0.2,0.1,null,null,null,null
   men_20_21, india, 1, 2, 3/11/2020,0.3,0.2,.0.1,null,null,null
   men_20_21, india, 2, 3, 4/11/2020,0.4,0.3,0.2,0.1,null,null
   men_20_21, india, 3, 4, 5/11/2020,0.5,0.4,0.3,0.2,0.1,null
   men_30_35, india, 4, 5, 6/11/2020,0.6,null,null,null,null,null
   men_30_35, india, 5, 6, 7/11/2020,0.7,0.6,null,null,null,null
   men_30_35, india, 6, 7, 8/11/2020,0.8,0.7,0.6,null,null,null
   men_30_35, india, 7, 8, 9/11/2020,0.9,0.8,0.7,0.6,null,null
   men_30_35, india, 8, 9, 10/11/2020,0.10,0.9,0.8,0.7,0.6,null
   men_30_35, india, 9, 10, 11/11/2020,0.11,0.9,0.8,0.7,0.6
   men_30_35, india, 10, 11, 12/11/2020,0.12,0.11,0.10,0.9,0.8,0.7

说明

当 gap>=date_nor 时,对于每一行选择评级列的最后 5[这里 n=5] 值,如果没有行满足条件更新,则所有最后 5 个值的值为空值。如果满足条件的行数小于 5,则将对应的最后一个值更新为 null。

我像下面这样加载了输入数据并进行了分组,但不知道如何继续。**

df=spark.read.csv("D:\\input\file1.csv")
df=df.groupby(["demo","origin"])

感谢任何帮助。谢谢。

【问题讨论】:

    标签: python-3.x apache-spark pyspark apache-spark-sql


    【解决方案1】:

    使用窗口函数我能够解决问题

    df.createOrReplaceTempView('df')
    df=spark.sql("select *, lag(ratings, 1, null) over (partition by demo, origin order by date_nor) as 1_last_rating,lag(ratings, 2, null) over (partition by demo, origin order by date_nor) as 2_last_rating,lag(ratings, 3, null) over (partition by demo, origin order by date_nor) as 3_last_rating,lag(ratings, 4, null) over (partition by demo, origin order by date_nor) as 4_last_rating,lag(ratings, 5, null) over (partition by demo, origin order by date_nor) as 5_last_rating from df")
    df.show()
    
    +---------+------+----+--------+-----------+-------+-------------+-------------+-------------+-------------+-------------+
    |     demo|origin| gap|date_nor|       date|ratings|1_last_rating|2_last_rating|3_last_rating|4_last_rating|5_last_rating|
    +---------+------+----+--------+-----------+-------+-------------+-------------+-------------+-------------+-------------+
    |men_20_21| india|-1.0|     0.0| 01/11/2020|    0.1|         null|         null|         null|         null|         null|
    |men_20_21| india| 0.0|     1.0| 02/11/2020|    0.2|          0.1|         null|         null|         null|         null|
    |men_20_21| india| 1.0|     2.0| 03/11/2020|    0.3|          0.2|          0.1|         null|         null|         null|
    |men_20_21| india| 2.0|     3.0| 04/11/2020|    0.4|          0.3|          0.2|          0.1|         null|         null|
    |men_20_21| india| 3.0|     4.0| 05/11/2020|    0.5|          0.4|          0.3|          0.2|          0.1|         null|
    |men_30_35| india| 4.0|     5.0| 06/11/2020|    0.6|         null|         null|         null|         null|         null|
    |men_30_35| india| 5.0|     6.0| 07/11/2020|    0.7|          0.6|         null|         null|         null|         null|
    |men_30_35| india| 6.0|     7.0| 08/11/2020|    0.8|          0.7|          0.6|         null|         null|         null|
    |men_30_35| india| 7.0|     8.0| 09/11/2020|    0.9|          0.8|          0.7|          0.6|         null|         null|
    |men_30_35| india| 8.0|     9.0| 10/11/2020|    0.1|          0.9|          0.8|          0.7|          0.6|         null|
    |men_30_35| india| 9.0|    10.0| 11/11/2020|   0.11|          0.1|          0.9|          0.8|          0.7|          0.6|
    |men_30_35| india|10.0|    11.0| 12/11/2020|   0.12|         0.11|          0.1|          0.9|          0.8|          0.7|
    +---------+------+----+--------+-----------+-------+-------------+-------------+-------------+-------------+-------------+
    

    【讨论】:

      【解决方案2】:

      您似乎希望根据某些业务条件创建一堆动态列。

      也许您可以利用“withColumn”功能,例如

      df.withColumn('last_rating_1', regex_replace(col('date_nor')...)

      【讨论】:

        猜你喜欢
        • 2022-01-12
        • 2022-12-21
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2020-03-07
        • 1970-01-01
        • 1970-01-01
        • 2012-07-07
        相关资源
        最近更新 更多