spark sql 窗口和补全有哪些功能
·
Spark SQL 窗口函数和补全功能是数据分析和处理的重要工具,以下是主要功能介绍:
一、窗口函数功能
1. 聚合类窗口函数
-- 基本窗口 SUM(col) OVER(PARTITION BY dept ORDER BY salary ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AVG(col) OVER(PARTITION BY category) COUNT(*) OVER(PARTITION BY user_id) MIN/MAX(col) OVER(PARTITION BY group)
2. 排名类窗口函数
ROW_NUMBER() OVER(PARTITION BY dept ORDER BY salary DESC) -- 唯一序号 RANK() OVER(PARTITION BY subject ORDER BY score DESC) -- 并列排名留空位 DENSE_RANK() OVER(PARTITION BY subject ORDER BY score) -- 并列排名不留空位 PERCENT_RANK() -- 百分比排名 NTILE(n) OVER(ORDER BY col) -- 数据分桶
3. 分析类窗口函数
LAG(col, n, default) OVER(ORDER BY time) -- 向前取第n行 LEAD(col, n, default) OVER(ORDER BY time) -- 向后取第n行 FIRST_VALUE(col) OVER(PARTITION BY id ORDER BY time) -- 窗口第一个值 LAST_VALUE(col) OVER(PARTITION BY id ORDER BY time ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING) -- 窗口最后一个值 NTH_VALUE(col, n) OVER(...) -- 第n个值
4. 统计类窗口函数
CUME_DIST() OVER(ORDER BY col) -- 累积分布 STDDEV(col) OVER(PARTITION BY ...) -- 标准差 VAR_POP(col) OVER(...) -- 总体方差
二、窗口定义选项
分区与排序
OVER( PARTITION BY col1, col2 -- 按列分组 ORDER BY time_col DESC -- 组内排序 [窗口框架] -- 定义行范围 )
窗口框架(Frame)
-- 行范围 ROWS BETWEEN 1 PRECEDING AND 1 FOLLOWING ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ROWS BETWEEN CURRENT ROW AND UNBOUNDED FOLLOWING -- 值范围(需要ORDER BY) RANGE BETWEEN INTERVAL 7 DAYS PRECEDING AND CURRENT ROW RANGE BETWEEN 100 PRECEDING AND 200 FOLLOWING
三、数据补全功能
1. 时间序列补全
-- 使用时间生成函数
SELECT
explode(sequence(
date_trunc('day', min_time),
date_trunc('day', max_time),
interval 1 day
)) as date
FROM source_table
-- 结合窗口函数进行插值
SELECT
time,
LAST_VALUE(value IGNORE NULLS) OVER(
ORDER BY time
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
) as interpolated_value
FROM table
2. 缺失值填充
-- 向前/向后填充
SELECT
id,
time,
COALESCE(
value,
LAG(value IGNORE NULLS) OVER(PARTITION BY id ORDER BY time),
LEAD(value IGNORE NULLS) OVER(PARTITION BY id ORDER BY time)
) as filled_value
FROM table
-- 使用均值填充
SELECT
id,
value,
COALESCE(value, AVG(value) OVER(PARTITION BY group)) as imputed_value
FROM table
3. 生成连续序列
-- 生成数字序列
SELECT explode(sequence(1, 10)) as num
-- 生成日期序列并左连接补全
WITH date_series AS (
SELECT explode(sequence(
'2024-01-01'::date,
'2024-01-31'::date,
interval 1 day
)) as date
)
SELECT
ds.date,
COALESCE(t.sales, 0) as sales
FROM date_series ds
LEFT JOIN sales_table t ON ds.date = t.sale_date
4. 区间补全
-- 使用CROSS JOIN生成所有组合
WITH all_categories AS (SELECT DISTINCT category FROM products),
all_dates AS (SELECT DISTINCT date FROM sales)
SELECT
c.category,
d.date,
COALESCE(s.amount, 0) as amount
FROM all_categories c
CROSS JOIN all_dates d
LEFT JOIN sales s ON c.category = s.category AND d.date = s.date
四、实用技巧组合
完整的时间序列分析
WITH filled_data AS (
SELECT
date,
user_id,
-- 用前一天数据填充当天缺失
COALESCE(
daily_active,
LAG(daily_active) OVER(PARTITION BY user_id ORDER BY date)
) as active_status,
-- 计算7天滚动平均
AVG(daily_active) OVER(
PARTITION BY user_id
ORDER BY date
ROWS BETWEEN 6 PRECEDING AND CURRENT ROW
) as rolling_avg_7d
FROM user_activity
)
会话标识与补全
SELECT
user_id,
timestamp,
SUM(session_start) OVER(
PARTITION BY user_id
ORDER BY timestamp
) as session_id
FROM (
SELECT *,
CASE
WHEN timestamp - LAG(timestamp) OVER(PARTITION BY user_id ORDER BY timestamp)
> interval 30 minutes
THEN 1
ELSE 0
END as session_start
FROM user_events
)
五、性能优化建议
-
合理使用分区:避免数据倾斜,分区键选择要均匀
-
框架范围优化:尽量使用
ROWS而非RANGE,前者性能更好 -
结合物化视图:对常用窗口计算结果进行物化
-
索引利用:确保ORDER BY列有合适的数据分布
这些功能组合使用可以解决复杂的时间序列分析、数据质量修复、报表生成等场景。
更多推荐




所有评论(0)