【问题标题】:Pyspark: how to create a new column and match the column's value condition with row valuePyspark:如何创建新列并将列的值条件与行值匹配
【发布时间】:2020-09-25 16:48:16
【问题描述】:
date        hos time topwait
19/9/2020   KHW 11:00   5
19/9/2020   CCM 11:00   6
19/9/2020   HHJ 11:00   7
19/9/2020   KHW 12:00   1
19/9/2020   CCM 12:00   2
19/9/2020   HHJ 12:00   4
22/9/2020   KHW 11:00   9
22/9/2020   CCM 11:00   9
22/9/2020   HHJ 11:00   9
22/9/2020   KHW 11:00   4
22/9/2020   CCM 11:00   3
22/9/2020   HHJ 11:00   2
.
.
.

我已获得如上所示的数据,我想添加一些新列以查看当天给定时间段的“topwait”。如下表所示。

date       hos     time topwait 11:00 12:00
19/9/2020   KHW    11:00    5   5   1
19/9/2020   CCM    11:00    6   6   2
19/9/2020   HHJ    11:00    7   7   4
19/9/2020   KHW    12:00    1   5   1
19/9/2020   CCM    12:00    2   6   2
19/9/2020   HHJ    12:00    4   7   4
22/9/2020   KHW    11:00    9   9   4
22/9/2020   CCM    11:00    9   9   3
22/9/2020   HHJ    11:00    9   9   2
22/9/2020   KHW    12:00    4   9   4
22/9/2020   CCM    12:00    3   9   3
22/9/2020   HHJ    12:00    2   9   2
.
.
.

我正在使用 pyspark 并尝试了下面的代码。但只有当时隙符合该列的条件,而其他的会显示为空时,它才能增加价值。

tt = df.withColumn('11:00',when((df.time == '11:00'), df.topwait))

date        hos    time topwait 11:00 12:00
19/9/2020   KHW    11:00    5   5     null
19/9/2020   CCM    11:00    6   6     null
19/9/2020   HHJ    11:00    7   7     null
19/9/2020   KHW    12:00    1   null    1
19/9/2020   CCM    12:00    2   null    2
19/9/2020   HHJ    12:00    4   null    4
22/9/2020   KHW    11:00    9   9     null
22/9/2020   CCM    11:00    9   9     null
22/9/2020   HHJ    11:00    9   9     null
22/9/2020   KHW    12:00    4   null    4
22/9/2020   CCM    12:00    3   null    3
22/9/2020   HHJ    12:00    2   null    2
.
.
.

我想它还需要以行的日期为条件,但我不知道如何以行值为条件。是否可以通过使用 pyspark 来实现?谢谢!

【问题讨论】:

    标签: python pyspark


    【解决方案1】:

    Pivot 和 join 也可以。

    df.join(df.groupBy('date', 'hos').pivot('time').agg(first('topwait').alias('topwait')), ['date', 'hos'], 'left') \
      .show()
    
    +---------+---+-----+-------+-----+-----+
    |     date|hos| time|topwait|11:00|12:00|
    +---------+---+-----+-------+-----+-----+
    |19/9/2020|KHW|11:00|      5|    5|    1|
    |19/9/2020|CCM|11:00|      6|    6|    2|
    |19/9/2020|HHJ|11:00|      7|    7|    4|
    |19/9/2020|KHW|12:00|      1|    5|    1|
    |19/9/2020|CCM|12:00|      2|    6|    2|
    |19/9/2020|HHJ|12:00|      4|    7|    4|
    |22/9/2020|KHW|11:00|      9|    9|    4|
    |22/9/2020|CCM|11:00|      9|    9|    3|
    |22/9/2020|HHJ|11:00|      9|    9|    2|
    |22/9/2020|KHW|12:00|      4|    9|    4|
    |22/9/2020|CCM|12:00|      3|    9|    3|
    |22/9/2020|HHJ|12:00|      2|    9|    2|
    +---------+---+-----+-------+-----+-----+
    

    【讨论】:

      【解决方案2】:

      基本上这不是一个简单的匹配。通过匹配它需要前一小时或后一小时的数据,为此我们需要date, hos 上的分区/组并按time 排序。 lag 下一小时,lead 前一小时。

      见下文-

      from pyspark.sql.window import Window
      from pyspark.sql import functions as f
      
      spark_df = sqlContext.createDataFrame([['19/9/2020','KHW','11:00','5'],
      ['19/9/2020','CCM','11:00','6'],
      ['19/9/2020','HHJ','11:00','7'],
      ['19/9/2020','KHW','12:00','1'],
      ['19/9/2020','CCM','12:00','2'],
      ['19/9/2020','HHJ','12:00','4'],
      ['22/9/2020','KHW','11:00','9'],
      ['22/9/2020','CCM','11:00','9'],
      ['22/9/2020','HHJ','11:00','9'],
      ['22/9/2020','KHW','12:00','4'],
      ['22/9/2020','CCM','12:00','3'],
      ['22/9/2020','HHJ','12:00','2']], ['date','hos', 'time', 'topwait'])
      
      spark_df.withColumn('11:00',
                          f.when(
                              spark_df.time == '11:00', 
                              spark_df.topwait
                          ).otherwise(
                              f.lag(spark_df.topwait).over(Window.partitionBy("date", "hos").orderBy("time")))
                          )\
      .withColumn('12:00',
                          f.when(
                              spark_df.time == '12:00', 
                              spark_df.topwait
                          ).otherwise(
                              f.lead(spark_df.topwait).over(Window.partitionBy("date", "hos").orderBy("time")))
                          )\
      .sort(["date", "time"]).show()   
      
      +---------+---+-----+-------+-----+-----+
      |     date|hos| time|topwait|11:00|12:00|
      +---------+---+-----+-------+-----+-----+
      |19/9/2020|KHW|11:00|      5|    5|    1|
      |19/9/2020|CCM|11:00|      6|    6|    2|
      |19/9/2020|HHJ|11:00|      7|    7|    4|
      |19/9/2020|KHW|12:00|      1|    5|    1|
      |19/9/2020|CCM|12:00|      2|    6|    2|
      |19/9/2020|HHJ|12:00|      4|    7|    4|
      |22/9/2020|HHJ|11:00|      9|    9|    2|
      |22/9/2020|KHW|11:00|      9|    9|    4|
      |22/9/2020|CCM|11:00|      9|    9|    3|
      |22/9/2020|CCM|12:00|      3|    9|    3|
      |22/9/2020|KHW|12:00|      4|    9|    4|
      |22/9/2020|HHJ|12:00|      2|    9|    2|
      +---------+---+-----+-------+-----+-----+
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2020-09-17
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2020-09-26
        • 2014-08-29
        相关资源
        最近更新 更多