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
)

五、性能优化建议

  1. 合理使用分区:避免数据倾斜,分区键选择要均匀

  2. 框架范围优化:尽量使用ROWS而非RANGE,前者性能更好

  3. 结合物化视图:对常用窗口计算结果进行物化

  4. 索引利用:确保ORDER BY列有合适的数据分布

这些功能组合使用可以解决复杂的时间序列分析、数据质量修复、报表生成等场景。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐