【问题标题】:Aggregating data at feature and time both在特征和时间上聚合数据
【发布时间】:2020-05-11 13:07:43
【问题描述】:

我有一个间隔 10 分钟的 pyspark 数据帧,我如何在一个分类特征和 2 小时的时间聚合它,然后计算其他两个特征的平均值和第三个特征的第一个值

我的示例数据在 pyspark 中如下所示。我想在 'ind' 和 'date' 2 小时的时间分组,然后计算 'sal' 的平均值和 'imp' 的第一个值

from pyspark import SparkContext
from pyspark.sql import SQLContext

sc = SparkContext.getOrCreate()
sqlContext = SQLContext(sc)

 a = sqlContext.createDataFrame([["Anand", "2020-02-01 16:00:00", 12, "ba"], 
                            ["Anand", "2020-02-01 16:10:00", 14,"sa"], 
                            ["Carl", "2020-02-01 16:00:00", 16,"da"], 
                            ["Carl", "2020-02-01 16:10:00", 12,"ga"],
                            ["Eric", "2020-02-01 16:o0:00", 24, "sa"]], ['ind', "date","sal","imp"])
a.show()

|  ind|               date|sal|imp|
+-----+-------------------+---+---+
|Anand|2020-02-01 16:00:00| 12| ba|
|Anand|2020-02-01 16:10:00| 14| sa|
| Carl|2020-02-01 16:00:00| 16| da|
| Carl|2020-02-01 16:10:00| 12| ga|
| Eric|2020-02-01 16:o0:00| 24| sa|

我不知道如何在 groupby Pyspark 中混合类别特征和时间(2 小时)。我知道如何在熊猫中做到这一点。但是我的真实数据是巨大的。有什么建议吗?

【问题讨论】:

  • 这是一个标准的 Spark SQL 问题,与 machine-learningpandasscikit-learn 无关 - 请不要向无关标签发送垃圾邮件(已删除)。
  • 你能在问题中添加预期的输出吗?

标签: pyspark apache-spark-sql


【解决方案1】:

为了更好地生成查询,我用几个时间示例扩展了您的案例

a = spark.createDataFrame([["Anand", "2020-02-01 16:00:00", 12, "ba"], 
                            ["Anand", "2020-02-01 16:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 17:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 18:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 19:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 20:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 21:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 22:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 23:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 00:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 01:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 02:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 03:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 04:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 05:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 06:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 07:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 08:10:00", 14,"sa"],
                           ["Anand", "2020-02-01 09:10:00", 14,"sa"],
                            ["Carl", "2020-02-01 16:00:00", 16,"da"], 
                            ["Carl", "2020-02-01 16:10:00", 12,"ga"],
                            ["Eric", "2020-02-01 16:00:00", 24, "sa"]], ['ind', "date","sal","imp"])

newa=a.withColumn('EveryTwoHour',f.when(f.hour(f.col('date').cast(t.TimestampType()))%2==0,
                                   f.hour(f.col('date').cast(t.TimestampType()))).otherwise(
                                   f.hour(f.col('date').cast(t.TimestampType()))+1))

newa.groupBy('ind','EveryTwoHour').agg(f.avg('sal'),f.first('imp')).orderBy('ind','EveryTwoHour').show()

+-----+------------+--------+-----------------+
|  ind|EveryTwoHour|avg(sal)|first(imp, false)|
+-----+------------+--------+-----------------+
|Anand|           0|    14.0|               sa|
|Anand|           2|    14.0|               sa|
|Anand|           4|    14.0|               sa|
|Anand|           6|    14.0|               sa|
|Anand|           8|    14.0|               sa|
|Anand|          10|    14.0|               sa|
|Anand|          16|    13.0|               ba|
|Anand|          18|    14.0|               sa|
|Anand|          20|    14.0|               sa|
|Anand|          22|    14.0|               sa|
|Anand|          24|    14.0|               sa|
| Carl|          16|    14.0|               da|
| Eric|          16|    24.0|               sa|
+-----+------------+--------+-----------------+

有多种方法可以做到,这只是其中之一。

为了每两小时执行一次聚合,我们为每个偶数小时创建一个新列,然后对其进行聚合。

a.withColumn('EveryTwoHour',f.when(f.hour(f.col('date').cast(t.TimestampType()))%2==0,
                                   f.hour(f.col('date').cast(t.TimestampType()))).otherwise(
    f.hour(f.col('date').cast(t.TimestampType()))+1)).show()

+-----+-------------------+---+---+------------+
|  ind|               date|sal|imp|EveryTwoHour|
+-----+-------------------+---+---+------------+
|Anand|2020-02-01 16:00:00| 12| ba|          16|
|Anand|2020-02-01 16:10:00| 14| sa|          16|
|Anand|2020-02-01 17:10:00| 14| sa|          18|
|Anand|2020-02-01 18:10:00| 14| sa|          18|
|Anand|2020-02-01 19:10:00| 14| sa|          20|
|Anand|2020-02-01 20:10:00| 14| sa|          20|
|Anand|2020-02-01 21:10:00| 14| sa|          22|
|Anand|2020-02-01 22:10:00| 14| sa|          22|
|Anand|2020-02-01 23:10:00| 14| sa|          24|
|Anand|2020-02-01 00:10:00| 14| sa|           0|
|Anand|2020-02-01 01:10:00| 14| sa|           2|
|Anand|2020-02-01 02:10:00| 14| sa|           2|
|Anand|2020-02-01 03:10:00| 14| sa|           4|
|Anand|2020-02-01 04:10:00| 14| sa|           4|
|Anand|2020-02-01 05:10:00| 14| sa|           6|
|Anand|2020-02-01 06:10:00| 14| sa|           6|
|Anand|2020-02-01 07:10:00| 14| sa|           8|
|Anand|2020-02-01 08:10:00| 14| sa|           8|
|Anand|2020-02-01 09:10:00| 14| sa|          10|
| Carl|2020-02-01 16:00:00| 16| da|          16|
+-----+-------------------+---+---+------------+

所以在这里,如果我正在获取小时并且如果它是偶数而不是没有变化并且如果小时是奇数,我将向它添加 1 以使其成为下一个偶数小时。

我希望它能解释逻辑。

【讨论】:

  • 嗨,Shubham,这真的很有帮助。您想添加一些关于创建两小时列的说明。它会帮助我更好
  • @ManuSharma 更新了答案
  • 嗨,我如何在日期和每两个小时重新采样相同的 pyspark 数据帧
猜你喜欢
  • 1970-01-01
  • 2021-04-02
  • 2017-02-21
  • 2014-01-26
  • 1970-01-01
  • 2021-12-20
  • 2021-05-30
  • 2020-01-20
  • 1970-01-01
相关资源
最近更新 更多