Spark SQL中开窗函数详解

简介: 笔记

row_number()开窗函数: 其实就是给每个分组的数据,按照其排序的顺序,打上一个分组内的行号,相当于groupTopN,在实际应用中非常广泛。

+--------+-------+------+
|deptName|   name|salary|
+--------+-------+------+
|  dept-1|Michael|  3000|
|  dept-2|   Andy|  5000|
|  dept-1|   Alex|  4500|
|  dept-2| Justin|  6700|
|  dept-2| Cherry|  3400|
|  dept-1|   Jack|  5500|
|  dept-2|   Jone| 12000|
|  dept-1|   Lucy|  8000|
|  dept-2|   LiLi|  7600|
|  dept-2|   Pony|  4200|
+--------+-------+------+

需求分析:对上面数据表按照deptName分组,并按照salary降序排序,取出每个deptName组前两名。

数据源:

{"deptName":"dept-1", "name":"Michael", "salary":3000}
{"deptName":"dept-2", "name":"Andy", "salary":5000}
{"deptName":"dept-1", "name":"Alex", "salary":4500}
{"deptName":"dept-2", "name":"Justin", "salary":6700}
{"deptName":"dept-2", "name":"Cherry", "salary":3400}
{"deptName":"dept-1", "name":"Jack", "salary":5500}
{"deptName":"dept-2", "name":"Jone", "salary":12000}
{"deptName":"dept-1", "name":"Lucy", "salary":8000}
{"deptName":"dept-2", "name":"LiLi", "salary":7600}
{"deptName":"dept-2", "name":"Pony", "salary":4200}

初始化SparkSession

package com.kfk.spark.common
import org.apache.spark.sql.SparkSession
/**
 * @author : 蔡政洁
 * @email :caizhengjie888@icloud.com
 * @date : 2020/12/2
 * @time : 10:02 下午
 */
object CommSparkSessionScala {
    def getSparkSession(): SparkSession ={
        val spark = SparkSession
                .builder
                .appName("CommSparkSessionScala")
                .master("local")
                .config("spark.sql.warehouse.dir", "/Users/caizhengjie/Document/spark/spark-warehouse")
                .getOrCreate
        return spark
    }
}

实现开窗函数

package com.kfk.spark.sql
import com.kfk.spark.common.{Comm, CommSparkSessionScala}
/**
 * @author : 蔡政洁
 * @email :caizhengjie888@icloud.com
 * @date : 2020/12/8
 * @time : 12:22 下午
 */
object WindowFunctionScala {
    def main(args: Array[String]): Unit = {
        val spark = CommSparkSessionScala.getSparkSession()
        val userPath = Comm.fileDirPath + "users.json"
        spark.read.json(userPath).show()
        /**
         * +--------+-------+------+
         * |deptName|   name|salary|
         * +--------+-------+------+
         * |  dept-1|Michael|  3000|
         * |  dept-2|   Andy|  5000|
         * |  dept-1|   Alex|  4500|
         * |  dept-2| Justin|  6700|
         * |  dept-2| Cherry|  3400|
         * |  dept-1|   Jack|  5500|
         * |  dept-2|   Jone| 12000|
         * |  dept-1|   Lucy|  8000|
         * |  dept-2|   LiLi|  7600|
         * |  dept-2|   Pony|  4200|
         * +--------+-------+------+
         */
        spark.read.json(userPath).createOrReplaceTempView("user")
        // 实现开窗函数:所谓开窗函数就是分组求TopN
        spark.sql("select deptName,name,salary,rank from" +
                "(select deptName,name,salary,row_number() OVER (PARTITION BY deptName order by salary desc) rank from user) tempUser " +
                "where rank <=2").show()
        /**
         * +--------+----+------+----+
         * |deptName|name|salary|rank|
         * +--------+----+------+----+
         * |  dept-1|Lucy|  8000|   1|
         * |  dept-1|Jack|  5500|   2|
         * |  dept-2|Jone| 12000|   1|
         * |  dept-2|LiLi|  7600|   2|
         * +--------+----+------+----+
         */
        // 实现分组排序
        spark.sql("select * from user order by deptName,salary desc").show()
        /**
         * +--------+-------+------+
         * |deptName|   name|salary|
         * +--------+-------+------+
         * |  dept-1|   Lucy|  8000|
         * |  dept-1|   Jack|  5500|
         * |  dept-1|   Alex|  4500|
         * |  dept-1|Michael|  3000|
         * |  dept-2|   Jone| 12000|
         * |  dept-2|   LiLi|  7600|
         * |  dept-2| Justin|  6700|
         * |  dept-2|   Andy|  5000|
         * |  dept-2|   Pony|  4200|
         * |  dept-2| Cherry|  3400|
         * +--------+-------+------+
         */
    }
}


相关文章
|
SQL JSON 分布式计算
Spark SQL架构及高级用法
Spark SQL基于Catalyst优化器与Tungsten引擎,提供高效的数据处理能力。其架构涵盖SQL解析、逻辑计划优化、物理计划生成及分布式执行,支持复杂数据类型、窗口函数与多样化聚合操作,结合自适应查询与代码生成技术,实现高性能大数据分析。
914 2
|
SQL 人工智能 数据挖掘
如何在`score`表中正确使用`COUNT`和`AVG`函数?SQL聚合函数COUNT与AVG使用指南
本文三桥君通过score表实例解析SQL聚合函数COUNT和AVG的常见用法。详解COUNT(studentNo)、COUNT(score)、COUNT()的区别,以及AVG函数对数值/字符型字段的不同处理,特别指出AVG()是无效语法。实战部分提供6个典型查询案例及结果,包含创建表、插入数据的完整SQL代码。产品专家三桥君强调正确理解函数特性(如空值处理、字段类型限制)对数据分析的重要性,帮助开发者避免常见误区,提升查询效率。
578 0
|
SQL 分布式计算 资源调度
Dataphin功能Tips系列(48)-如何根据Hive SQL/Spark SQL的任务优先级指定YARN资源队列
如何根据Hive SQL/Spark SQL的任务优先级指定YARN资源队列
557 4
|
SQL 分布式计算 Java
Spark SQL向量化执行引擎框架Gluten-Velox在AArch64使能和优化
本文摘自 Arm China的工程师顾煜祺关于“在 Arm 平台上使用 Native 算子库加速 Spark”的分享,主要内容包括以下四个部分: 1.技术背景 2.算子库构成 3.算子操作优化 4.未来工作
2318 0
|
SQL 数据库 数据库管理
数据库SQL函数应用技巧与方法
在数据库管理中,SQL函数是处理和分析数据的强大工具
|
SQL Oracle 关系型数据库
SQL优化-使用联合索引和函数索引
在一次例行巡检中,发现一条使用 `to_char` 函数将日期转换为字符串的 SQL 语句 CPU 利用率很高。为了优化该语句,首先分析了 where 条件中各列的选择性,并创建了不同类型的索引,包括普通索引、函数索引和虚拟列索引。通过对比不同索引的执行计划,最终确定了使用复合索引(包含函数表达式)能够显著降低查询成本,提高执行效率。
449 3
|
SQL 数据库 索引
SQL中COUNT函数结合条件使用的技巧与方法
在SQL查询中,COUNT函数是一个非常常用的聚合函数,用于计算表中满足特定条件的记录数
3175 5
|
SQL JSON 分布式计算
【赵渝强老师】Spark SQL的数据模型:DataFrame
本文介绍了在Spark SQL中创建DataFrame的三种方法。首先,通过定义case class来创建表结构,然后将CSV文件读入RDD并关联Schema生成DataFrame。其次,使用StructType定义表结构,同样将CSV文件读入RDD并转换为Row对象后创建DataFrame。最后,直接加载带有格式的数据文件(如JSON),通过读取文件内容直接创建DataFrame。每种方法都包含详细的代码示例和解释。
528 0
|
SQL 分布式计算 数据库
【大数据技术Spark】Spark SQL操作Dataframe、读写MySQL、Hive数据库实战(附源码)
【大数据技术Spark】Spark SQL操作Dataframe、读写MySQL、Hive数据库实战(附源码)
1067 0
|
SQL 分布式计算 大数据
【大数据技术Hadoop+Spark】Spark SQL、DataFrame、Dataset的讲解及操作演示(图文解释)
【大数据技术Hadoop+Spark】Spark SQL、DataFrame、Dataset的讲解及操作演示(图文解释)
760 0