为了更好地生成查询,我用几个时间示例扩展了您的案例
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 以使其成为下一个偶数小时。
我希望它能解释逻辑。