
摘要离线数仓里“每个类目销售额前三的商品”“用户最长连续登录天数”环比增长率这类需求用 groupBy 往往要写多层自连接。窗口函数一条 SQL 就能解决。这篇文章梳理 over 子句的窗口规范用六个真实场景把 ROW_NUMBER、LAG、SUM/AVG OVER 的写法和坑讲清楚。关键词窗口函数, OVER, ROW_NUMBER, LAG, 累计求和, 连续登录, TopN一、先搞清楚窗口函数解决什么问题做订单分析时常有这种需求既要每个订单的明细又要在每一行上算这个用户到目前为止累计消费了多少。用 groupBy 会折叠行拿到的是每个用户一个汇总值明细就丢了想同时保留明细和汇总只能把聚合结果 join 回原表。窗口函数的价值就在这里——它在一行上算一个基于窗口一组相关行的聚合值但不折叠任何一行。语法骨架FUNC(col)OVER(PARTITIONBYkey-- 分区按 key 分组组内独立计算ORDERBYsort_col-- 排序决定窗口内行的先后ROWS/RANGEBETWEEN...-- 帧边界滑动窗口范围可选)FUNC 分三类排名ROW_NUMBER / RANK / DENSE_RANK、偏移LAG / LEAD、聚合SUM / AVG / COUNT / MIN / MAX 加 OVER。二、两个最容易被搞混的点2.1 ROWS 与 RANGE 的区别这是窗口函数里最容易写错的地方。差别在于滑动窗口怎么界定-- ROWS按物理行数偏移ROWSBETWEEN1PRECEDINGANDCURRENTROW-- 当前行往上数 1 行就是窗口不管值是否相等-- RANGE按排序键的值偏移RANGEBETWEEN1PRECEDINGANDCURRENTROW-- ORDER BY 列的值在 [当前值-1, 当前值] 的所有行一句话ROWS 数的是行RANGE 数的是值。RANGE 要求 ORDER BY 的列是数值或日期否则没法做值减 1这种运算。2.2 默认帧的坑不写 ROWS/RANGE 时Spark 的默认帧取决于有没有 ORDER BY-- 有 ORDER BY默认是累计-- RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROWSUM(amount)OVER(PARTITIONBYuidORDERBYdt)-- 逐行累计-- 无 ORDER BY默认是全区SUM(amount)OVER(PARTITIONBYuid)-- 每行都是该 uid 的总和这一点在写 SUM OVER 累计时不用显式声明但心里要清楚默认行为否则改个排序方向结果就变了。三、六个生产场景下面是数仓里真正高频的六类需求每个都给出能直接跑的 SQL 和踩坑点。场景一每个类目销售额 Top 3 商品需求电商报表里最经典的 TopN每个 category 取 sales 最高的 3 个商品。SELECTcategory,product_id,salesFROM(SELECTcategory,product_id,sales,ROW_NUMBER()OVER(PARTITIONBYcategoryORDERBYsalesDESC)ASrnFROMproduct_sales)tWHERErn3;两个坑并列怎么处理ROW_NUMBER 会给并列值随机分配一个唯一序号可能漏掉本来并列第三的商品。如果业务要求并列都保留用DENSE_RANK()1,1,2,3 不跳号把rn 3换成rank 3。分区倾斜。某个 category 有上亿行时这一个分区会拖垮整个 job。应对办法见第四节。场景二用户最长连续登录天数需求留存分析里统计每个用户连续登录的最长天数。这是经典的 Gaps and Islands 问题。核心思路连续日期的特点是日期递增 1行号也递增 1所以日期 - 行号是个常量能标识一段连续区间。SELECTuid,MAX(days)ASmax_consecutive_daysFROM(SELECTuid,grp,COUNT(*)ASdaysFROM(SELECTuid,login_date,DATE_SUB(login_date,ROW_NUMBER()OVER(PARTITIONBYuidORDERBYlogin_date))ASgrpFROM(SELECTDISTINCTuid,login_dateFROMuser_login)dedup)groupedGROUPBYuid,grp)aggGROUPBYuid;两个坑必须先去重。一个用户同一天登录多次不去重的话日期-行号会被打乱连续判断直接失效。上面用SELECT DISTINCT uid, login_date先干掉重复。DATE_SUB(date, n)在 SparkSQL 里是date - n 天第二个参数是 ROW_NUMBER 返回的整数类型对得上才能减。场景三环比 / 同比增长率需求每个部门月度销售额的环比增长率。SELECTdept,month,sales,LAG(sales,1)OVER(PARTITIONBYdeptORDERBYmonth)ASprev_sales,ROUND((sales-LAG(sales,1)OVER(PARTITIONBYdeptORDERBYmonth))*100.0/LAG(sales,1)OVER(PARTITIONBYdeptORDERBYmonth),2)ASgrowth_pctFROMdept_sales_monthly;两个坑LAG 的默认值。每个分区第一行没有上一行LAG(sales, 1)返回 NULLNULL 参与除法会让整行 growth 变成 NULL。给第三个参数一个默认值LAG(sales, 1, 0)或者用CASE WHEN显式处理第一行。环比是LAG(..., 1)同比去年同期是LAG(..., 12)按月粒度。别把偏移量写反。场景四累计求和Running Total需求每个用户按时间累计消费金额常用于等级、额度计算。SELECTuid,order_date,amount,SUM(amount)OVER(PARTITIONBYuidORDERBYorder_dateROWSBETWEENUNBOUNDEDPRECEDINGANDCURRENTROW)AScumulative_amountFROMorders;这里显式写ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW是防御性写法——虽然它是默认帧但显式声明能让读代码的人一眼看清这是累计不用去回忆默认行为。场景五每个用户取最新一条记录去重需求用户表有历史变更流水要每个 uid 的最新一条。SELECTuid,order_id,amount,tsFROM(SELECTuid,order_id,amount,ts,ROW_NUMBER()OVER(PARTITIONBYuidORDERBYtsDESC)ASrnFROMorder_log)tWHERErn1;和groupBy(uid) max(ts)相比窗口函数能保留整行不需要再 join 一次把最新记录的其它字段取回来。这是它在这个场景下的核心价值。场景六7 日移动平均需求销量趋势图里平滑毛刺算 7 日移动平均。SELECTdt,sales,AVG(sales)OVER(ORDERBYdtROWSBETWEEN6PRECEDINGANDCURRENTROW)ASma7FROMdaily_sales;注意这里没写 PARTITION BY是全量单序列的移动平均。如果按商品分别算补上PARTITION BY product_id。窗口范围是当前行 前 6 行 7 天所以是6 PRECEDING别写成7 PRECEDING。四、性能与踩坑窗口函数会触发 Shuffle按 PARTITION BY 的 key 重分布代价不低。生产上留意几点分区键选高基数列。PARTITION BY 的列基数太小比如性别、布尔值数据会挤进少数几个分区退化成单点计算。先过滤再开窗。WHERE 能下推就下推先把数据量压下去再算窗口别在大表上直接开窗再过滤。避免全量 ROW_NUMBER 再去重。如果只是要每类最大的一个groupBy max更便宜ROW_NUMBER 全量排序的代价要高一个量级。并列名次要谨慎。业务口径是取前 N 名还是取前 N 行决定了用 DENSE_RANK 还是 ROW_NUMBER这在数仓里常出数据不一致的事故。五、总结窗口函数的本质是聚合但不折叠行需要明细 汇总同时出现的场景优先考虑。排名用 ROW_NUMBER唯一还是 DENSE_RANK并列保留取决于业务口径这是最容易埋数据口径坑的地方。写 SUM/AVG OVER 时心里明确默认帧是累计需要滑动窗口就显式写 ROWS/RANGE。连续登录、去重保留最新这两个场景窗口函数比 groupBy join 更简洁也更省一次 Shuffle。作者大数据技术实践者博客blog.starzy.cnGitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践