【问题标题】:pyspark rolling window timeframepyspark 滚动窗口时间范围
【发布时间】:2020-08-13 09:07:22
【问题描述】:

我正在尝试实现一个按 source_ip 分组的 30 分钟时间范围的滚动窗口。我们的想法是获取每个 source_ip 的平均值。不确定这是正确的方法。我遇到的问题是 ip 192.168.1.3 似乎平均超过 30 分钟的窗口,因为数据包 25 是几天后。

df = sqlContext.createDataFrame([('192.168.1.1', 17, "2017-03-10T15:27:18+00:00"),
                        ('192.168.1.2', 1, "2017-03-15T12:27:18+00:00"),
                        ('192.168.1.2', 2, "2017-03-15T12:28:18+00:00"),
                        ('192.168.1.2', 3, "2017-03-15T12:29:18+00:00"),
                        ('192.168.1.3', 4, "2017-03-15T12:28:18+00:00"),
                        ('192.168.1.3', 5, "2017-03-15T12:29:18+00:00"),
                        ('192.168.1.3', 25, "2017-03-18T11:27:18+00:00")],
                        ["source_ip","packets", "timestampGMT"])

w = (Window()
     .partitionBy("source_ip")
     .orderBy(F.col("timestampGMT").cast('long'))
     .rangeBetween(-1800, 0))

df = df.withColumn('rolling_average', F.avg("packets").over(w))

df.show(100,False)

这是我得到的结果。我希望前 2 个条目为 4.5,第三个条目为 25?

+-----------+-------+-------------------------+------------------+
|source_ip  |packets|timestampGMT             |rolling_average   |
+-----------+-------+-------------------------+------------------+
|192.168.1.3|4      |2017-03-15T12:28:18+00:00|11.333333333333334|
|192.168.1.3|5      |2017-03-15T12:29:18+00:00|11.333333333333334|
|192.168.1.3|25     |2017-03-18T11:27:18+00:00|11.333333333333334|
|192.168.1.2|1      |2017-03-15T12:27:18+00:00|2.0               |
|192.168.1.2|2      |2017-03-15T12:28:18+00:00|2.0               |
|192.168.1.2|3      |2017-03-15T12:29:18+00:00|2.0               |
|192.168.1.1|17     |2017-03-10T15:27:18+00:00|17.0              |
+-----------+-------+-------------------------+------------------+

【问题讨论】:

    标签: pyspark


    【解决方案1】:

    先把字符串改成时间戳,然后orderBy它。

    import pyspark.sql.functions as F
    from pyspark.sql import Window
    
    w = (Window()
         .partitionBy("source_ip")
         .orderBy(F.col("timestamp"))
         .rangeBetween(-1800, 0))
    
    df = df.withColumn("timestamp", F.unix_timestamp(F.to_timestamp("timestampGMT"))) \
        .withColumn('rolling_average', F.avg("packets").over(w))
    
    df.printSchema()
    df.show(100,False)
    
    
    root
     |-- source_ip: string (nullable = true)
     |-- packets: long (nullable = true)
     |-- timestampGMT: string (nullable = true)
     |-- timestamp: long (nullable = true)
     |-- rolling_average: double (nullable = true)
    
    +-----------+-------+-------------------------+----------+---------------+
    |source_ip  |packets|timestampGMT             |timestamp |rolling_average|
    +-----------+-------+-------------------------+----------+---------------+
    |192.168.1.2|1      |2017-03-15T12:27:18+00:00|1489580838|1.0            |
    |192.168.1.2|2      |2017-03-15T12:28:18+00:00|1489580898|1.5            |
    |192.168.1.2|3      |2017-03-15T12:29:18+00:00|1489580958|2.0            |
    |192.168.1.1|17     |2017-03-10T15:27:18+00:00|1489159638|17.0           |
    |192.168.1.3|4      |2017-03-15T12:28:18+00:00|1489580898|4.0            |
    |192.168.1.3|5      |2017-03-15T12:29:18+00:00|1489580958|4.5            |
    |192.168.1.3|25     |2017-03-18T11:27:18+00:00|1489836438|25.0           |
    +-----------+-------+-------------------------+----------+---------------+
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-03-13
      • 1970-01-01
      • 2021-09-11
      • 2021-07-21
      相关资源
      最近更新 更多