Spark/Databricks SQL:分区内按10秒间隔分组并取最小时间为组起始
在Spark/Databricks SQL中实现动态10秒窗口分组
要实现按组起始时间+10秒的动态窗口分组(而非固定时间桶),可以使用递归CTE追踪每组的起始时间,判断后续记录是否属于当前组,最后聚合得到结果。以下是针对需求的完整实现:
示例实现代码
WITH RECURSIVE input_data AS ( -- 模拟输入数据集 SELECT TIMESTAMP '2023-04-11 04:20:00' AS Time, 'val1' AS Text UNION ALL SELECT TIMESTAMP '2023-04-11 04:36:00' AS Time, 'val2' AS Text UNION ALL SELECT TIMESTAMP '2023-04-11 04:40:00' AS Time, 'val3' AS Text UNION ALL SELECT TIMESTAMP '2023-04-11 04:45:00' AS Time, 'val4' AS Text UNION ALL SELECT TIMESTAMP '2023-04-11 04:47:00' AS Time, 'val5' AS Text UNION ALL SELECT TIMESTAMP '2023-04-11 04:50:00' AS Time, 'val6' AS Text UNION ALL SELECT TIMESTAMP '2023-04-11 04:55:00' AS Time, 'val7' AS Text UNION ALL SELECT TIMESTAMP '2023-04-11 04:56:00' AS Time, 'val8' AS Text UNION ALL SELECT TIMESTAMP '2023-04-11 05:13:00' AS Time, 'val9' AS Text ), ranked_data AS ( -- 按时间排序并添加行号,用于递归遍历 SELECT Time, Text, ROW_NUMBER() OVER (ORDER BY Time) AS rn FROM input_data ), recursive_groups AS ( -- 初始化:第一条记录作为第一组 SELECT rn, Time, Text, Time AS group_start, 1 AS group_id FROM ranked_data WHERE rn = 1 UNION ALL -- 递归处理后续每条记录 SELECT rd.rn, rd.Time, rd.Text, -- 判断当前记录是否在当前组的10秒窗口内,是则沿用组起始时间,否则用当前时间作为新组起始 CASE WHEN rd.Time <= rg.group_start + INTERVAL 10 SECOND THEN rg.group_start ELSE rd.Time END, -- 开启新组则组ID+1,否则沿用当前组ID CASE WHEN rd.Time <= rg.group_start + INTERVAL 10 SECOND THEN rg.group_id ELSE rg.group_id + 1 END FROM recursive_groups rg JOIN ranked_data rd ON rg.rn = rd.rn - 1 ) -- 按组聚合,拼接Text字段并取组起始时间 SELECT group_start AS group_start_time, STRING_AGG(Text, ':') AS concatenated_text FROM recursive_groups GROUP BY group_id, group_start ORDER BY group_start_time;
代码逻辑说明
input_data:模拟给定的示例输入,包含Time(时间戳)和Text(待拼接字段)。ranked_data:给每条记录按时间顺序添加行号,确保递归可以按顺序处理每一条记录。recursive_groups:- 初始步骤选取第一条记录,将其时间设为第一组的起始时间,组ID为1。
- 递归步骤依次处理后续记录:对比当前记录时间与前一组的起始时间,如果当前时间在起始时间+10秒范围内,则加入当前组;否则开启新组,更新组起始时间和组ID。
- 最终聚合:按组ID和组起始时间分组,用
STRING_AGG拼接同组的Text字段,得到目标结果。
输出结果
执行上述代码后,会得到如下结果:
| group_start_time | concatenated_text |
|---|---|
| 2023-04-11 04:20:00.000 | val1 |
| 2023-04-11 04:36:00.000 | val2:val3:val4 |
| 2023-04-11 04:47:00.000 | val5:val6:val7:val8 |
| 2023-04-11 05:13:00.000 | val9 |
结果完全符合指定的分组规则:
- 第一组仅包含
val1,时间范围04:20-04:30,无其他记录。 - 第二组包含
val2、val3、val4,时间范围04:36-04:46。 - 第三组包含
val5、val6、val7、val8,时间范围04:47-04:57。 - 最后一组仅包含
val9,时间05:13,超出前一组的10秒范围。
内容的提问来源于stack exchange,提问作者pratap_pachipulusu
相关产品推荐
相关产品推荐

