【问题标题】:Pyspark Group and Order by Sum for Group Divide by partsPyspark Group 和 Order by Sum for Group 除以部分
【发布时间】:2020-04-23 11:45:47
【问题描述】:

我是 pyspark 的新手,我很困惑如何将一些数据按几列组合在一起,按另一列排序,然后为每个组添加一列,然后将其用作每一行的分母用于计算组成组的每一行中的权重的数据。

这是在 jupyterlab 中使用 pyspark3 笔记本完成的。没有办法解决这个问题。

这是数据的一个例子......

+-------+-----+-----------+------------+------+--------+
| ntwrk | zip | zip-ntwrk | event-date | hour | counts |
+-------+-----+-----------+------------+------+--------+
| A     | 1   | 1-A       | 2019-10-10 | 1    | 12362  |
| B     | 3   | 3-B       | 2019-10-10 | 1    | 100    |
| C     | 5   | 5-C       | 2019-10-10 | 1    | 17493  |
| B     | 3   | 3-B       | 2019-10-10 | 4    | 4873   |
| A     | 2   | 2-A       | 2019-10-11 | 1    | 28730  |
| C     | 6   | 6-C       | 2019-10-11 | 1    | 728    |
| C     | 5   | 5-C       | 2019-10-10 | 2    | 9827   |
| A     | 1   | 1-A       | 2019-10-10 | 9    | 13245  |
| B     | 4   | 4-B       | 2019-10-11 | 1    | 3765   |
+-------+-----+-----------+------------+------+--------+

我想按 ntrk、zipcode、zip-ntwrk、event-date 将其组合在一起,然后按 event-date desc 和 hour desc 对其进行排序。每个日期有 24 小时,因此对于每个 zip-ntwrk 组合,我希望按顺序查看日期和小时。像这样的...

+-------+-----+-----------+------------+------+--------+
| ntwrk | zip | zip-ntwrk | event-date | hour | counts |
+-------+-----+-----------+------------+------+--------+
| A     | 1   | 1-A       | 2019-10-10 | 1    | 12362  |
| A     | 1   | 1-A       | 2019-10-10 | 9    | 3765   |
| A     | 2   | 2-A       | 2019-10-11 | 1    | 28730  |
| B     | 3   | 3-B       | 2019-10-10 | 1    | 100    |
| B     | 3   | 3-B       | 2019-10-10 | 4    | 4873   |
| B     | 4   | 4-B       | 2019-10-11 | 1    | 3765   |
| C     | 5   | 5-C       | 2019-10-10 | 1    | 17493  |
| C     | 5   | 5-C       | 2019-10-10 | 2    | 9827   |
| C     | 6   | 6-C       | 2019-10-11 | 1    | 728    |
+-------+-----+-----------+------------+------+--------+

现在一切都井井有条,我需要运行一个计算来计算每小时计数与每天计数总和的比率。这将在分母中用于将每小时计数除以总数,以获得每小时计数与一天总数的比率。所以像这样......

+-------+-----+-----------+------------+------+--------+-------+
| ntwrk | zip | zip-ntwrk | event-date | hour | counts | total |
+-------+-----+-----------+------------+------+--------+-------+
| A     | 1   | 1-A       | 2019-10-10 | 1    | 12362  | 16127 |
| A     | 1   | 1-A       | 2019-10-10 | 9    | 3765   | 16127 |
| A     | 2   | 2-A       | 2019-10-11 | 1    | 28730  | 28730 |
| B     | 3   | 3-B       | 2019-10-10 | 1    | 100    | 4973  |
| B     | 3   | 3-B       | 2019-10-10 | 4    | 4873   | 4973  |
| B     | 4   | 4-B       | 2019-10-11 | 1    | 3765   | 3765  |
| C     | 5   | 5-C       | 2019-10-10 | 1    | 17493  | 27320 |
| C     | 5   | 5-C       | 2019-10-10 | 2    | 9827   | 27320 |
| C     | 6   | 6-C       | 2019-10-11 | 1    | 728    | 728   |
+-------+-----+-----------+------------+------+--------+-------+

现在我们有了分母,我们可以将每行的计数除以总数得到因子 counts/total=factor,这最终看起来像......

+-------+-----+-----------+------------+------+--------+-------+--------+
| ntwrk | zip | zip-ntwrk | event-date | hour | counts | total | factor |
+-------+-----+-----------+------------+------+--------+-------+--------+
| A     | 1   | 1-A       | 2019-10-10 | 1    | 12362  | 16127 | .766   |
| A     | 1   | 1-A       | 2019-10-10 | 9    | 3765   | 16127 | .233   |
| A     | 2   | 2-A       | 2019-10-11 | 1    | 28730  | 28730 | 1      |
| B     | 3   | 3-B       | 2019-10-10 | 1    | 100    | 4973  | .02    |
| B     | 3   | 3-B       | 2019-10-10 | 4    | 4873   | 4973  | .979   |
| B     | 4   | 4-B       | 2019-10-11 | 1    | 3765   | 3765  | 1      |
| C     | 5   | 5-C       | 2019-10-10 | 1    | 17493  | 27320 | .64    |
| C     | 5   | 5-C       | 2019-10-10 | 2    | 9827   | 27320 | .359   |
| C     | 6   | 6-C       | 2019-10-11 | 1    | 728    | 728   | 1      |
+-------+-----+-----------+------------+------+--------+-------+--------+

这就是我正在尝试做的事情,任何关于如何完成这件事的建议都将不胜感激。

谢谢

【问题讨论】:

    标签: python algorithm apache-spark pyspark


    【解决方案1】:

    使用window sum函数,然后通过ntwrk,zip对窗口分区进行sum

    • 最后我们将与counts/total分开。

    Example:

    from pyspark.sql.functions import *
    from pyspark.sql import Window
    w = Window.partitionBy("ntwrk","zip","event-date")
    
    df1.withColumn("total",sum(col("counts")).over(w).cast("int")).orderBy("ntwrk","zip","event-date","hour").\
    withColumn("factor",format_number(col("counts")/col("total"),3)).show()
    
    #+-----+---+---------+----------+----+------+-----+------+
    #|ntwrk|zip|zip-ntwrk|event-date|hour|counts|total|factor|
    #+-----+---+---------+----------+----+------+-----+------+
    #|    A|  1|      1-A|2019-10-10|   1| 12362|25607| 0.483|
    #|    A|  1|      1-A|2019-10-10|   9| 13245|25607| 0.517|#input 13245 not 3765
    #|    A|  2|      2-A|2019-10-11|   1| 28730|28730| 1.000|
    #|    B|  3|      3-B|2019-10-10|   1|   100| 4973| 0.020|
    #|    B|  3|      3-B|2019-10-10|   4|  4873| 4973| 0.980|
    #|    B|  4|      4-B|2019-10-11|   1|  3765| 3765| 1.000|
    #|    C|  5|      5-C|2019-10-10|   1| 17493|27320| 0.640|
    #|    C|  5|      5-C|2019-10-10|   2|  9827|27320| 0.360|
    #|    C|  6|      6-C|2019-10-11|   1|   728|  728| 1.000|
    #+-----+---+---------+----------+----+------+-----+------+
    

    【讨论】:

    • 这是一个很好的开始。它为每个 ntwrk 和 zip 进行了计算,但它没有在日期之前完成,所以我不得不添加一些片段。非常感谢。这是我在下面使用的代码...from pyspark.sql.functions import *from pyspark.sql import Windoww = Window.partitionBy("ntwrk","zip","event-date")factor_df = group_by_dataframe.withColumn("total",sum(col("counts")).over(w).` cast("int")).orderBy("ntwrk","zip","event-date","hour")。 `withColumn("factor",format_number(col("counts")/col("total"),3))factor_df.show(50)
    【解决方案2】:

    你一定是网状样条线

    【讨论】:

      【解决方案3】:

      Pyspark 适用于分布式架构,因此它可能不会保留顺序。因此,在显示记录之前,您应该始终按照您需要的方式对其进行排序。

      现在,您的重点是获取各个级别的记录百分比。您可以使用窗口功能实现相同的功能,按您想要的数据级别进行分区。

      喜欢: w = Window.partitionBy("ntwrk-zip", "hour") df =df.withColumn("hourly_recs", F.count().over(w))

      另外,你可以参考 YouTube 上的这个教程 - https://youtu.be/JEBd_4wWyj0

      【讨论】:

        猜你喜欢
        • 2017-02-06
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2021-12-24
        • 2020-01-19
        • 1970-01-01
        • 2011-06-28
        相关资源
        最近更新 更多