逸Y   仙X头像
关注
Spark SQL封面图

Spark SQL

Spark SQL介绍

Spark SQL是Apache Spark生态系统中用于处理结构化数据的核心模块。它的核心作用可以概括为三点:一是提供统一的数据访问接口,能用同样的SQL或API处理来自JSON、Parquet、Hive表等多种数据源的数据;二是内置了Catalyst优化器,能自动优化查询逻辑,大幅提升大数据查询的执行效率;三是与Hive生态高度兼容,能无缝集成现有数据仓库,并支持JDBC/ODBC等标准连接方式。

Spark SQL的发展史源于对性能的极致追求。其前身Shark(Hive on Spark)虽然通过将底层引擎替换为Spark提升了查询速度,但因架构臃肿而被放弃。2014年,Spark团队推出了完全重写的Spark SQL,并持续演进:1.3版本定义了DataFrame,1.6版本引入了强类型的Dataset API,3.0版本增加了自适应查询执行(AQE)等智能化特性,实现了统一、高效的结构化数据处理。

理解Spark on Hive与Hive on Spark的区别至关重要。Spark on Hive是由Spark主导、将Hive仅作为元数据源和数据存储的方案,性能更优,适合新项目;而Hive on Spark则是由Hive主导、仅将Spark作为底层执行引擎的方案,主要用于从原有Hive生态平滑迁移。当前的Spark SQL已不再依赖Hive的执行组件,但仍需依赖Hive Metastore来管理表元数据,这种模块化设计使其既保持了强大的功能,又具备灵活的部署能力。

Spark SQL编程

添加依赖

<dependency>
  <groupId>org.apache.spark</groupId>
  <artifactId>spark-sql_2.12</artifactId>
  <version>3.3.1</version>
</dependency>

构建环境对象

public static void main(String[] args) {
    //创建环境对象
    //直接new是不推荐的
    //        SparkContext sparkContext = new SparkContext(new SparkConf().setMaster("local[3]").setAppName("learn_test"));
    //        SparkSession sparkSession = new SparkSession(sparkContext);

    //推荐使用的是构建器模式
    SparkSession sparkSession = SparkSession.builder()
    .appName("learn_test")
    .master("local[2]")
    .getOrCreate();


    //释放资源
    sparkSession.close();
}
public class SparkSQL01_Env_01 {

    public static void main(String[] args) {
        //创建环境对象
        //直接new是不推荐的
//        SparkContext sparkContext = new SparkContext();
//        SparkSession sparkSession = new SparkSession(sparkContext);
        
        SparkConf sparkConf = new SparkConf().setMaster("local[3]").setAppName("learn_test");
     
        //推荐使用的是构建器模式
        SparkSession sparkSession = SparkSession.builder()
                .config(sparkConf)
                .getOrCreate();


        //释放资源
        sparkSession.close();
    }
}

快速开始

public class SparkSQL02_Model {

    public static void main(String[] args) {
        //创建环境对象
        //直接new是不推荐的
        //        SparkContext sparkContext = new SparkContext(new SparkConf().setMaster("local[3]").setAppName("learn_test"));
        //        SparkSession sparkSession = new SparkSession(sparkContext);

        //推荐使用的是构建器模式
        SparkSession sparkSession = SparkSession.builder()
        .appName("learn_test")
        .master("local[2]")
        .getOrCreate();

        //使用sparkSession读取文件
        //现在的数据模型变成了DataSet Spark之后统一使用DataSet
        Dataset<Row> ds = sparkSession.read().json("doc/user.json");

        //将数据转化成历临时的表结构进行查询
        ds.createOrReplaceTempView("user");


        String sql = "select * from user";
        //使用sql速去
        Dataset<Row> sqled = sparkSession.sql(sql);

        sqled.show();
        //释放资源
        sparkSession.close();
    }
}

DSL使用

DSL 是 Spark 提供的基于方法链的编程接口,让你可以用 Scala/Java 的语法来编写类似于 SQL 的查询,而不需要写字符串形式的 SQL 语句。

// SQL 方式(字符串拼接,容易出错,无编译时检查)
sparkSession.sql("SELECT name, age FROM user WHERE age > 18 ORDER BY age DESC");

// DSL 方式(类型安全,IDE 智能提示,链式调用)
ds.select("name", "age")
  .where(col("age").gt(18))
  .orderBy(col("age").desc())
  .show();

代码展示:

public class SparkSQL02_Model_01 {

    public static void main(String[] args) {
        //创建环境对象
        //直接new是不推荐的
//        SparkContext sparkContext = new SparkContext(new SparkConf().setMaster("local[3]").setAppName("learn_test"));
//        SparkSession sparkSession = new SparkSession(sparkContext);

        //推荐使用的是构建器模式
        SparkSession sparkSession = SparkSession.builder()
                .appName("learn_test")
                .master("local[2]")
                .getOrCreate();

        //使用sparkSession读取文件
        //现在的数据模型变成了DataSet Spark之后统一使用DataSet
        Dataset<Row> ds = sparkSession.read().json("doc/user.json");
        //使用User + Encoder显示具体类型
        //Dataset<User> userDataset = ds.as(Encoders.bean(User.class));
        //userDs.filter(u -> u.getAge() > 18).show();
        //使用dsl的方式访问
        Dataset<Row> select = ds.select("*");

        select.show();
        //释放资源
        sparkSession.close();
    }
}
class User implements Serializable {
    private Integer id;
    private String name;
    private double salary;
    private String city;
    private Integer age;
    private Boolean active;
}

ufd

Spark SQL中的UDF(User-Defined Function,用户自定义函数),是一种让开发者能够根据自身业务逻辑扩展Spark SQL功能的编程接口。当Spark提供的数百种内置函数(如 sumsplitrank 等)无法满足特定需求时,你可以通过UDF定义自己的处理逻辑,并像调用内置函数一样在SQL语句中使用它。

public class SparkSQL03_udf {

    public static void main(String[] args) {
        //创建环境对象
        //直接new是不推荐的
//        SparkContext sparkContext = new SparkContext(new SparkConf().setMaster("local[3]").setAppName("learn_test"));
//        SparkSession sparkSession = new SparkSession(sparkContext);

        //推荐使用的是构建器模式
        SparkSession sparkSession = SparkSession.builder()
                .appName("learn_test")
                .master("local[2]")
                .getOrCreate();

        //使用sparkSession读取文件
        //现在的数据模型变成了DataSet Spark之后统一使用DataSet
        Dataset<Row> ds = sparkSession.read().json("doc/user.json");

        //将数据转化成历临时的表结构进行查询
        ds.createOrReplaceTempView("user");

        //使用udf 自定义方法
        //todo : StringType$.MODULE$ 这里使用的是scala 的语法
        sparkSession.udf().register("perfixName",name->"name:"+name, StringType$.MODULE$);
        //这里是Java的方式
        sparkSession.udf().register("perfixName",name->"name:"+name, DataTypes.StringType);

        //在sql中使用
        String sql = "select perfixName(name) from user";
        //使用sql查询
        Dataset<Row> sqled = sparkSession.sql(sql);

        sqled.show();
        //释放资源
        sparkSession.close();
    }
}

udaf

UDAF(User-Defined Aggregation Function,用户自定义聚合函数) 是 Spark SQL 中另一类重要的用户自定义函数。它与 UDF 的核心区别在于:UDF 对“单行”数据进行处理并返回一个值,而 UDAF 对“多行”数据进行聚合处理并返回一个值

public static void main(String[] args) {
        //创建环境对象
        //直接new是不推荐的
        //SparkContext sparkContext = new SparkContext(new SparkConf().setMaster("local[3]").setAppName("learn_test"));
        //SparkSession sparkSession = new SparkSession(sparkContext);

        //推荐使用的是构建器模式
        SparkSession sparkSession = SparkSession.builder()
                .appName("learn_test")
                .master("local[2]")
                .getOrCreate();

        //使用sparkSession读取文件
        //现在的数据模型变成了DataSet Spark之后统一使用DataSet
        Dataset<Row> ds = sparkSession.read().json("doc/user.json");

        //将数据转化成历临时的表结构进行查询
        ds.createOrReplaceTempView("user");

        //使用udaf 自定义方法 计算年龄的平均值(这个是模拟数据)
        sparkSession.udf().register("getAvgAge", functions.udaf(new AvgUserAgeDuaf(), Encoders.INT()));


        //在sql中使用
        String sql = "select getAvgAge(age) from user";
        //使用sql查询
        Dataset<Row> sqled = sparkSession.sql(sql);

        sqled.show();
        //释放资源
        sparkSession.close();
    }
public class AvgUserAgeDuaf extends Aggregator<Integer, AvgAgerBuffer, Integer> {


    //初始化缓冲区
    @Override
    public AvgAgerBuffer zero() {
        return new AvgAgerBuffer(0,0);
    }

    //缓冲区的聚合操作
    @Override
    public AvgAgerBuffer reduce(AvgAgerBuffer buffer, Integer a) {
        buffer.setTotal(buffer.getTotal() + a);
        buffer.setCnt(buffer.getCnt() + 1);
        return buffer;
    }


    //合并缓冲区,他是分布式框架,所以需要进行缓冲区的合并
    @Override
    public AvgAgerBuffer merge(AvgAgerBuffer b1, AvgAgerBuffer b2) {
        b1.setTotal(b1.getTotal() + b2.getTotal());
        b1.setCnt(b1.getCnt() + b2.getCnt());
        return b1;
    }

    //计算最终结果
    @Override
    public Integer finish(AvgAgerBuffer reduction) {
        Integer ans = reduction.getTotal() / reduction.getCnt();
        return ans;
    }

    @Override
    public Encoder<AvgAgerBuffer> bufferEncoder() {
        return Encoders.bean(AvgAgerBuffer.class);
    }

    @Override
    public Encoder<Integer> outputEncoder() {
        return Encoders.INT();
    }
}
public class AvgAgerBuffer implements Serializable {
    //数值和
    private Integer total;

    //计数
    private Integer cnt;


    public Integer getTotal() {
        return total;
    }

    public void setTotal(Integer total) {
        this.total = total;
    }

    public Integer getCnt() {
        return cnt;
    }

    public void setCnt(Integer cnt) {
        this.cnt = cnt;
    }

    public AvgAgerBuffer() {
    }

    public AvgAgerBuffer(Integer total, Integer cnt) {
        this.total = total;
        this.cnt = cnt;
    }
}

数据的加载和保存

在 Spark SQL 中,最常用的三种文件类型是:JSON、CSV、Parquet。这也是生产环境中使用频率最高的三种数据源。

对比维度

JSON

CSV

Parquet

JDBC

Hive

格式类型

半结构化文本

纯文本表格

列式二进制

关系型数据库

数据仓库元数据

人类可读

✅ 是

✅ 是

❌ 否

❌ 否

❌ 否

自带 Schema

✅ 是

❌ 需推断

✅ 是

✅ 是(通过 JDBC)

✅ 是(Metastore)

支持嵌套

✅ 支持

❌ 不支持

✅ 支持

❌ 不支持

✅ 支持(取决于存储格式)

压缩率

中等

中等

极高

取决于 DB

取决于存储格式

读取速度

较慢

中等

最快

中等(受网络影响)

较快

数据规模

GB~TB

GB~TB

PB 级

GB~TB

PB 级

典型场景

API 日志、接口数据

Excel 导出、数据交换

数据湖/数仓存储

业务库导入导出

已有 Hadoop 数仓

CSV

public class SparkSQL05_Source_csv {

    public static void main(String[] args) {
        //创建环境对象
        //直接new是不推荐的
//        SparkContext sparkContext = new SparkContext(new SparkConf().setMaster("local[3]").setAppName("learn_test"));
//        SparkSession sparkSession = new SparkSession(sparkContext);

        //推荐使用的是构建器模式
        SparkSession sparkSession = SparkSession.builder()
                .appName("learn_test")
                .master("local[2]")
                .getOrCreate();

        //读取csv文件
        Dataset<Row> ds = sparkSession
                .read()
                .option("header",true)
                .csv("doc/user.csv");

        ds.createOrReplaceTempView("user");
        //这里现在不是实际的字段名字,看着语义不清晰

//        String sql = "select avg(_c2) from user";
        //todo: 语义清晰一点
        String sql = "select avg(age) from user";

        //清晰表达方式,在read中添加配置,告诉spark我们的数据是又表头的

        Dataset<Row> sqled = sparkSession.sql(sql);

        sqled.show();

        //释放资源
        sparkSession.close();
    }
}

Json

public static void main(String[] args) {
        //创建环境对象
        //直接new是不推荐的
//        SparkContext sparkContext = new SparkContext(new SparkConf().setMaster("local[3]").setAppName("learn_test"));
//        SparkSession sparkSession = new SparkSession(sparkContext);

        //推荐使用的是构建器模式
        SparkSession sparkSession = SparkSession.builder()
                .appName("learn_test")
                .master("local[2]")
                .getOrCreate();

        //读取csv文件
        Dataset<Row> ds = sparkSession
                .read()
                .json("doc/user.json");

        ds.createOrReplaceTempView("user");
        //这里现在不是实际的字段名字,看着语义不清晰

//        String sql = "select avg(_c2) from user";
        //todo: 语义清晰一点
        String sql = "select * from user";

        //清晰表达方式,在read中添加配置,告诉spark我们的数据是又表头的

        Dataset<Row> sqled = sparkSession.sql(sql);

        sqled.show();

        //释放资源
        sparkSession.close();
    }

Parquet

public static void main(String[] args) {

        //推荐使用的是构建器模式
        SparkSession sparkSession = SparkSession.builder()
                .appName("learn_test")
                .master("local[2]")
                .getOrCreate();

        //读取csv文件
        Dataset<Row> ds = sparkSession
                .read()
                .parquet("doc/user.parquet");

        ds.createOrReplaceTempView("user");
        //这里现在不是实际的字段名字,看着语义不清晰
        
        String sql = "select * from user";

        //清晰表达方式,在read中添加配置,告诉spark我们的数据是又表头的

        Dataset<Row> sqled = sparkSession.sql(sql);


        sqled.show();

        //释放资源
        sparkSession.close();
    }

JDBC

Spark SQL 可以通过 JDBC 连接器从 MySQL 数据库中读取数据。具体来说,用户可以使用 Spark 提供的 DataFrameReader 接口,并通过指定数据源类型为 jdbc 来实现与 MySQL 的连接。在连接过程中,需要提供 MySQL 的数据库地址(URL)、用户名、密码以及要读取的表名等信息。此外,还需要将 MySQL 的 JDBC 驱动程序(如 mysql-connector-java)添加到 Spark 应用程序的依赖中,以确保能够成功建立数据库连接。

public class SparkSQL05_Source_MySQL {


    private static String url = "jdbc:mysql://localhost:3306/yd?serverTimezone=Asia/Shanghai&useSSL=false&useUnicode=true&characterEncoding=UTF-8&allowPublicKeyRetrieval=true";

    private static String tableName = "user";

    private static Properties properties = new Properties();

    static {
        properties.setProperty("user", "******");
        properties.setProperty("password", "*******");
    }

    public static void main(String[] args) {

        //推荐使用的是构建器模式
        SparkSession sparkSession = SparkSession.builder()
                .appName("learn_test")
                .master("local[2]")
                .getOrCreate();

        //读取csv文件
        Dataset<Row> ds = sparkSession
                .read()
                .jdbc(url, tableName, properties);

        ds.createOrReplaceTempView("user");
        //这里现在不是实际的字段名字,看着语义不清晰

        String sql = "select * from user";

        //清晰表达方式,在read中添加配置,告诉spark我们的数据是又表头的

        Dataset<Row> sqled = sparkSession.sql(sql);


        Dataset<UserDB> userDataset = sqled.as(Encoders.bean(UserDB.class));

        //过滤出角色为admin的用户
        Dataset<UserDB> adminDataset = userDataset.filter(
                (FilterFunction<UserDB>) user -> "admin".equals(user.getRole()));

        adminDataset.show();

        //释放资源
        sparkSession.close();
    }
}
class UserDB {
    private Integer id;

    private String userName;

    private String fullName;

    private String password;

    private String email;

    private String role;

    private String mobile;

    private Boolean status;

    private LocalDateTime createTime;

    private LocalDateTime updateTime;

    public UserDB() {
    }

    public Integer getId() {
        return id;
    }

    public void setId(Integer id) {
        this.id = id;
    }

    public String getUserName() {
        return userName;
    }

    public void setUserName(String userName) {
        this.userName = userName;
    }

    public String getFullName() {
        return fullName;
    }

    public void setFullName(String fullName) {
        this.fullName = fullName;
    }

    public String getPassword() {
        return password;
    }

    public void setPassword(String password) {
        this.password = password;
    }

    public String getEmail() {
        return email;
    }

    public void setEmail(String email) {
        this.email = email;
    }

    public String getRole() {
        return role;
    }

    public void setRole(String role) {
        this.role = role;
    }

    public String getMobile() {
        return mobile;
    }

    public void setMobile(String mobile) {
        this.mobile = mobile;
    }

    public Boolean getStatus() {
        return status;
    }

    public void setStatus(Boolean status) {
        this.status = status;
    }

    public LocalDateTime getCreateTime() {
        return createTime;
    }

    public void setCreateTime(LocalDateTime createTime) {
        this.createTime = createTime;
    }

    public LocalDateTime getUpdateTime() {
        return updateTime;
    }

    public void setUpdateTime(LocalDateTime updateTime) {
        this.updateTime = updateTime;
    }
}

Hive

Park SQL 将 Hive 作为数据源
Park SQL 是一种用于处理大规模数据的分布式计算框架,它能够通过结构化查询语言(SQL)对数据进行高效的数据提取、转换和分析。在实际应用中,Park SQL 常常与 Hive 集成,将 Hive 作为其主要的数据源之一。
Hive 是基于 Hadoop 的一个数据仓库工具,它可以将结构化的数据文件映射为一张数据库表,并支持类 SQL 查询功能,通过 HiveQL(Hive Query Language)将 SQL 语句转换为 MapReduce 任务运行。由于 Hive 提供了类似传统数据库的查询方式,且能够处理海量数据,因此被广泛用于大数据分析领域。
当 Park SQL 将 Hive 作为数据源时,可以通过配置相应的连接参数(如 Hive 的元数据存储地址、HiveServer2 的连接信息等),直接读取 Hive 中存储的数据表。这种集成方式通常依赖于 Hive 的元数据管理功能,使得 Park SQL 能够自动识别 Hive 表的结构(如字段名、字段类型等),从而简化数据读取和查询的流程。
此外,Park SQL 对 Hive 的支持通常通过特定的连接器(Connector)或模块(如 Spark Hive Integration)实现。这种连接器不仅能够读取 Hive 表的数据,还支持将查询结果写入 Hive 表中,从而实现双向的数据交互。
总的来说,Park SQL 将 Hive 作为数据源的集成方式,使得用户能够充分利用 Hive 的数据存储能力和 Park SQL 的高性能计算能力,为大数据分析提供了强大的支持。这种组合在企业级数据仓库和实时分析场景中具有广泛的应用价值。

转载自 CSDN-专业IT技术社区

原文链接:https://blog.csdn.net/hdk5855/article/details/163998070

文章来源转载

评论

赞0

评论列表

微信小程序
QQ小程序

关于作者

点赞数:0
关注数:0
粉丝:0
文章:0
关注标签:0
加入于:--