无法理解PySpark中window时间窗口函数使用后的输出值
代码输出结果
你执行上述代码后,Databricks中会返回3行结果,结构如下(具体时间显示会匹配你环境的时区设置,以下为东八区环境下的参考输出):
| window | sum(val) |
|---|---|
| {"start": "0000-01-01T00:00:00.000+08:00", "end": "0000-01-01T00:00:55.000+08:00"} | 1 |
| {"start": "1970-01-01T19:01:15.000+08:00", "end": "1970-01-01T19:02:10.000+08:00"} | 1 |
| {"start": "1970-01-01T19:02:10.000+08:00", "end": "1970-01-01T19:03:05.000+08:00"} | 1 |
window函数运行原理
pyspark.sql.functions.window是Spark为时间序列数据设计的分组函数,用来将时间数据按照指定规则划分到不同的时间窗口中,后续可基于窗口做聚合计算,核心逻辑如下:
- 基础模式判断:如果不传
slideDuration参数,默认其值等于windowDuration,此时生成的是滚动窗口,窗口之间无重叠、无间隔,每条数据只会属于一个窗口。 - 参数说明:
timeColumn:要用来划分窗口的时间字段,必须为Timestamp类型,传入字符串会触发隐式类型转换,转换失败会报错。windowDuration:单个窗口的时间长度,支持seconds、minutes、hours、days等时间单位。slideDuration:窗口滑动的间隔,小于窗口长度时生成有重叠的滑动窗口,一条数据可能属于多个窗口;大于窗口长度时生成有间隔的跳跃窗口,部分数据不会被划入任何窗口。startTime:窗口对齐的偏移量,默认以Unix epoch时间(1970-01-01 00:00:00 UTC)为对齐基准。
- 窗口计算规则:
- 将时间列的每个值转换为相对于epoch的时间戳(早于1970年的时间对应负数时间戳)。
- 用时间戳减去
startTime偏移量后,对slideDuration向下取整,得到窗口的起始时间戳。 - 窗口结束时间 = 起始时间戳 +
windowDuration。 - 所有落入同一个[起始时间, 结束时间)区间的数据会被分到同一个组,执行后续的聚合操作。
小提示:生产环境建议手动用to_timestamp函数将时间字符串转为Timestamp类型再传入window,避免隐式转换带来的格式、时区异常。
内容的提问来源于stack exchange,提问作者Chitransh Mathur
相关产品推荐
相关产品推荐

