Spark SQL 中的窗口函数是一种强大的分析工具,允许在数据集的逻辑窗口上执行计算。让我详细解释其实现方式和常见函数。

窗口函数的基本概念

窗口定义结构

function_name() OVER (
    [PARTITION BY partition_expression, ...]
    [ORDER BY sort_expression [ASC|DESC], ...]
    [frame_clause]
)

核心组件说明

  1. PARTITION BY: 定义数据分区(类似GROUP BY)
  2. ORDER BY: 定义分区内排序规则
  3. Frame Clause: 定义窗口框架范围

常见窗口函数分类

1. 排名函数

-- ROW_NUMBER(): 行号(无重复)
ROW_NUMBER() OVER (PARTITION BY department ORDER BY salary DESC)

-- RANK(): 排名(有间隔)
RANK() OVER (PARTITION BY department ORDER BY salary DESC)

-- DENSE_RANK(): 密集排名(无间隔)
DENSE_RANK() OVER (PARTITION BY department ORDER BY salary DESC)

-- PERCENT_RANK(): 百分比排名
PERCENT_RANK() OVER (ORDER BY score)

-- CUME_DIST(): 累积分布
CUME_DIST() OVER (ORDER BY score)

2. 分析函数

-- LAG/LEAD: 获取前一行或后一行数据
LAG(salary, 1) OVER (PARTITION BY dept ORDER BY hire_date)
LEAD(salary, 1) OVER (PARTITION BY dept ORDER BY hire_date)

-- FIRST_VALUE/LAST_VALUE: 第一个/最后一个值
FIRST_VALUE(name) OVER (PARTITION BY dept ORDER BY salary DESC)
LAST_VALUE(name) OVER (PARTITION BY dept ORDER BY salary DESC ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING)

3. 聚合函数作为窗口函数

-- SUM/COUNT/AVG/MIN/MAX等聚合函数
SUM(salary) OVER (PARTITION BY department ORDER BY hire_date)
COUNT(*) OVER (PARTITION BY department)
AVG(salary) OVER (PARTITION BY department)

实际应用示例

示例数据准备

CREATE TEMPORARY VIEW employee_data AS
SELECT * FROM VALUES
(1, 'Alice', 'Engineering', 80000, '2020-01-15'),
(2, 'Bob', 'Engineering', 75000, '2019-03-20'),
(3, 'Charlie', 'Sales', 60000, '2021-05-10'),
(4, 'David', 'Sales', 65000, '2020-08-25'),
(5, 'Eve', 'Marketing', 70000, '2019-11-30')
AS t(id, name, department, salary, hire_date);

典型应用场景

1. 计算部门内薪资排名
SELECT 
    name,
    department,
    salary,
    ROW_NUMBER() OVER (PARTITION BY department ORDER BY salary DESC) as salary_rank,
    RANK() OVER (PARTITION BY department ORDER BY salary DESC) as salary_rank_with_ties,
    DENSE_RANK() OVER (PARTITION BY department ORDER BY salary DESC) as dense_salary_rank
FROM employee_data;
2. 计算累计薪资总额
SELECT 
    name,
    department,
    salary,
    hire_date,
    SUM(salary) OVER (PARTITION BY department ORDER BY hire_date) as cumulative_salary
FROM employee_data
ORDER BY department, hire_date;
3. 比较员工与部门平均薪资
SELECT 
    name,
    department,
    salary,
    AVG(salary) OVER (PARTITION BY department) as avg_dept_salary,
    salary - AVG(salary) OVER (PARTITION BY department) as salary_diff_from_avg
FROM employee_data;
4. 查找前后记录
SELECT 
    name,
    department,
    salary,
    hire_date,
    LAG(name, 1) OVER (PARTITION BY department ORDER BY hire_date) as previous_hire,
    LEAD(name, 1) OVER (PARTITION BY department ORDER BY hire_date) as next_hire,
    LAG(salary, 1) OVER (PARTITION BY department ORDER BY hire_date) as previous_salary
FROM employee_data
ORDER BY department, hire_date;

DataFrame API 实现方式

Scala 示例

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._

val windowSpec = Window.partitionBy("department").orderBy(desc("salary"))

val result = df.withColumn("rank", row_number().over(windowSpec))
              .withColumn("dense_rank", dense_rank().over(windowSpec))
              .withColumn("cumulative_sum", sum("salary").over(
                Window.partitionBy("department")
                      .orderBy("hire_date")
                      .rowsBetween(Window.unboundedPreceding, Window.currentRow)))

Python 示例

from pyspark.sql import Window
from pyspark.sql.functions import *

windowSpec = Window.partitionBy("department").orderBy(desc("salary"))

result = df.withColumn("rank", row_number().over(windowSpec)) \
           .withColumn("dense_rank", dense_rank().over(windowSpec)) \
           .withColumn("cumulative_sum", sum("salary").over(
             Window.partitionBy("department")
                   .orderBy("hire_date")
                   .rowsBetween(Window.unboundedPreceding, Window.currentRow)))

窗口框架类型

ROWS vs RANGE

-- ROWS: 基于物理行数
ROWS BETWEEN 2 PRECEDING AND 2 FOLLOWING

-- RANGE: 基于逻辑值范围
RANGE BETWEEN 100 PRECEDING AND 100 FOLLOWING

常用框架定义

-- 当前行到最后一行
ROWS BETWEEN CURRENT ROW AND UNBOUNDED FOLLOWING

-- 从第一行到当前行
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW

-- 当前行前后各一行
ROWS BETWEEN 1 PRECEDING AND 1 FOLLOWING

性能优化建议

1. 分区策略

-- 合理使用PARTITION BY减少数据量
-- 避免过多小分区或过少大分区

2. 排序优化

-- 只在必要时使用ORDER BY
-- 避免不必要的多字段排序

3. 缓存策略

-- 对频繁使用的窗口结果进行缓存
CACHE TABLE windowed_results

高级窗口函数

NTILE 函数

-- 将数据分为N个桶
NTILE(4) OVER (ORDER BY salary) as quartile

NTH_VALUE 函数

-- 获取第N个值
NTH_VALUE(name, 2) OVER (PARTITION BY department ORDER BY salary DESC) as second_highest_paid

窗口函数是 Spark SQL 中非常强大的功能,能够高效地处理复杂的分析需求,如排名、趋势分析、比较分析等,在数据分析和报表生成中具有重要价值。

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐