Spark Scala实现基于时间戳的30分钟滑动窗口连续行分组
数据表
| space_id | template | frequency | day | timestamp |
|---|---|---|---|---|
| 321d8 | temp | 15 | 2023-02-22T00:00:00+05:30 | 2023-02-22T09:00:00+05:30 |
| 321d8 | temp | 15 | 2023-02-22T00:00:00+05:30 | 2023-02-22T09:15:00+05:30 |
| 321d8 | temp | 15 | 2023-02-22T00:00:00+05:30 | 2023-02-22T09:30:00+05:30 |
| 321d8 | temp | 15 | 2023-02-22T00:00:00+05:30 | 2023-02-22T09:45:00+05:30 |
| 321d8 | temp | 15 | 2023-02-22T00:00:00+05:30 | 2023-02-22T10:00:00+05:30 |
字段说明
space_id:唯一标识IDtemplate:数据类型(如温度、湿度、CO₂等)frequency:传感器数据上报频率day:日期字段timestamp:数据上报的时间戳
需求说明
当前已实现非重叠的30分钟批次分组,但实际需要基于timestamp的30分钟滑动窗口分组:以每条数据的timestamp作为窗口起始点,包含该时间往后30分钟内的所有行,最终形成如下重叠分组:
- 第1、2、3行(窗口起始09:00,覆盖至09:30)
- 第2、3、4行(窗口起始09:15,覆盖至09:45)
- 第3、4、5行(窗口起始09:30,覆盖至10:00)
实现方案(SQL)
方法1:自连接匹配时间范围
通过自连接关联同空间内、时间戳在当前行往后30分钟内的数据,得到每个滑动窗口的分组:
SELECT t1.timestamp AS window_start, t2.space_id, t2.template, t2.frequency, t2.day, t2.timestamp AS data_timestamp FROM sensor_data t1 JOIN sensor_data t2 ON t2.timestamp >= t1.timestamp AND t2.timestamp <= t1.timestamp + INTERVAL '30 minutes' AND t1.space_id = t2.space_id ORDER BY t1.timestamp, t2.timestamp;
方法2:时间范围窗口函数(支持的数据库:PostgreSQL 11+、BigQuery等)
利用支持时间范围的窗口函数,直接定义滑动窗口并聚合数据:
SELECT timestamp AS window_start, space_id, template, frequency, day, ARRAY_AGG(timestamp) OVER ( PARTITION BY space_id ORDER BY timestamp RANGE BETWEEN CURRENT ROW AND INTERVAL '30 minutes' FOLLOWING ) AS window_data_timestamps FROM sensor_data;
内容的提问来源于stack exchange,提问作者Rohit M
相关产品推荐
相关产品推荐

