【发布时间】:2019-10-08 02:12:15
【问题描述】:
使用两个预定义的边界在 Spark SQL 中指定窗口间隔的正确方法是什么?
我正在尝试在“3 小时前到 2 小时前”的窗口中汇总表中的值。
当我运行这个查询时:
select *, sum(value) over (
partition by a, b
order by cast(time_value as timestamp)
range between interval 2 hours preceding and current row
) as sum_value
from my_temp_table;
这行得通。我得到了我期望的结果,即落在 2 小时滚动窗口内的值的总和。
现在,我需要的是让滚动窗口不绑定到当前行,而是考虑 3 小时前和 2 小时前之间的行。 我试过了:
select *, sum(value) over (
partition by a, b
order by cast(time_value as timestamp)
range between interval 3 hours preceding and 2 hours preceding
) as sum_value
from my_temp_table;
但我收到extraneous input 'hours' expecting {'PRECEDING', 'FOLLOWING'} 错误。
我也试过了:
select *, sum(value) over (
partition by a, b
order by cast(time_value as timestamp)
range between interval 3 hours preceding and interval 2 hours preceding
) as sum_value
from my_temp_table;
然后我得到不同的错误scala.MatchError: CalendarIntervalType (of class org.apache.spark.sql.types.CalendarIntervalType$)
我尝试的第三个选项是:
select *, sum(value) over (
partition by a, b
order by cast(time_value as timestamp)
range between interval 3 hours preceding and 2 preceding
) as sum_value
from my_temp_table;
它并没有像我们预期的那样工作:cannot resolve 'RANGE BETWEEN interval 3 hours PRECEDING AND 2 PRECEDING' due to data type mismatch
我很难找到间隔类型的文档,因为this link 说得不够多,其他信息也有点半生不熟。至少我发现了什么。
【问题讨论】:
-
AFAIK 范围间隔目前在 SparkSQL 中无法正常工作,只有基于行数的间隔才是稳健的。请参阅此 JIRA 票 issues.apache.org/jira/browse/SPARK-25842 。弃用也标记为 Scala API github.com/apache/spark/blob/v2.4.3/sql/core/src/main/scala/org/…
-
我明白了。好的,所以我要去寻找另一种方法。
标签: apache-spark apache-spark-sql window-functions