You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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;

代码逻辑说明

  1. input_data:模拟给定的示例输入,包含Time(时间戳)和Text(待拼接字段)。
  2. ranked_data:给每条记录按时间顺序添加行号,确保递归可以按顺序处理每一条记录。
  3. recursive_groups:
    • 初始步骤选取第一条记录,将其时间设为第一组的起始时间,组ID为1。
    • 递归步骤依次处理后续记录:对比当前记录时间与前一组的起始时间,如果当前时间在起始时间+10秒范围内,则加入当前组;否则开启新组,更新组起始时间和组ID。
  4. 最终聚合:按组ID和组起始时间分组,用STRING_AGG拼接同组的Text字段,得到目标结果。

输出结果

执行上述代码后,会得到如下结果:

group_start_timeconcatenated_text
2023-04-11 04:20:00.000val1
2023-04-11 04:36:00.000val2:val3:val4
2023-04-11 04:47:00.000val5:val6:val7:val8
2023-04-11 05:13:00.000val9

结果完全符合指定的分组规则:

  • 第一组仅包含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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.25 06:30:09