窗口函数条件聚合:缺口与岛屿问题实现用户意见周期统计
用户意见持续周期统计实现方案
需求背景
- 源数据字段:
user_id、question、opinion(取值1/0)、last_modified(意见提交时间),存在同一用户同一问题的连续重复意见记录 - 统计目标:输出每个用户每个问题下连续相同意见的时间周期,包含字段
user_id、question、opinion、from(周期起始时间)、until(周期结束时间,最新状态的该字段为NULL) - 现存问题:基础行偏移窗口函数未合并连续重复意见,会生成冗余错误条目
可行实现方案
纯SQL实现(适用于ADF映射数据流的SQL源查询、或ADB/Spark SQL等计算引擎)
核心思路用间隙与岛屿(Gaps and Islands)算法标记连续相同意见的分组,再分组聚合得到周期:
WITH ordered_data AS ( -- 按用户、问题分组,按修改时间升序排序 SELECT user_id, question, opinion, last_modified, -- 标记和上一条记录意见不同的行,作为新分组的起点 LAG(opinion, 1, -1) OVER (PARTITION BY user_id, question ORDER BY last_modified) AS prev_opinion FROM your_source_table ), group_mark AS ( SELECT *, -- 累计求和生成连续相同意见的分组ID SUM(CASE WHEN opinion != prev_opinion THEN 1 ELSE 0 END) OVER (PARTITION BY user_id, question ORDER BY last_modified) AS group_id FROM ordered_data ) -- 分组聚合得到每个连续意见段的起止时间 SELECT user_id, question, opinion, MIN(last_modified) AS `from`, -- 最后一个分组的结束时间置为NULL CASE WHEN ROW_NUMBER() OVER (PARTITION BY user_id, question ORDER BY group_id DESC) = 1 THEN NULL ELSE MAX(last_modified) END AS `until` FROM group_mark GROUP BY user_id, question, opinion, group_id ORDER BY user_id, question, `from`
ADF可视化数据流实现(无需写SQL)
按以下顺序添加变换组件即可:
- 源变换:接入你的原始数据表
- 排序变换:排序条件设置为
user_id升序、question升序、last_modified升序 - 窗口变换:
- 分区依据:
user_id、question - 排序依据:
last_modified升序 - 添加列
prev_opinion,表达式为lag(opinion, 1, -1)
- 分区依据:
- 派生列变换:添加列
is_new_group,表达式为iif(opinion != prev_opinion, 1, 0) - 窗口变换:
- 分区依据:
user_id、question - 排序依据:
last_modified升序 - 添加列
group_id,表达式为sum(is_new_group),窗口帧选择「无界起始到当前行」
- 分区依据:
- 聚合变换:
- 分组依据:
user_id、question、opinion、group_id - 聚合列
from:聚合规则为min(last_modified) - 聚合列
max_time:聚合规则为max(last_modified) - 聚合列
group_rn:聚合规则为row_number(sort(group_id desc))
- 分组依据:
- 派生列变换:添加列
until,表达式为iif(group_rn == 1, null(), max_time) - 选择变换:保留输出字段
user_id、question、opinion、from、until - 接收器变换:写入目标表即可
内容的提问来源于stack exchange,提问作者d1sh4
相关产品推荐
相关产品推荐

