在 Spark SQL 中,如何实现窗口函数操作?常见的窗口函数有哪些?
·
Spark SQL 中的窗口函数是一种强大的分析工具,允许在数据集的逻辑窗口上执行计算。让我详细解释其实现方式和常见函数。
窗口函数的基本概念
窗口定义结构
function_name() OVER (
[PARTITION BY partition_expression, ...]
[ORDER BY sort_expression [ASC|DESC], ...]
[frame_clause]
)
核心组件说明
- PARTITION BY: 定义数据分区(类似GROUP BY)
- ORDER BY: 定义分区内排序规则
- 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 中非常强大的功能,能够高效地处理复杂的分析需求,如排名、趋势分析、比较分析等,在数据分析和报表生成中具有重要价值。
更多推荐


所有评论(0)