sparksql20问-1
1、什么是 Spark SQL?它的主要功能是什么?
spark sql 是 spark 用于处理结构化数据的模块,提供了 DataFrame API 的编程抽象,并且可以无缝整合进 spark 的 其他组件,允许用户执行 sql 查询,读取数据、转换数据,并将数据保存到不同的存储系统中。
主要功能:支持 ANSI SQL 标准,能够进行复杂查询分析;能够与多种数据源整合,包括 Hive、Cassandra、HBase、JSON、CSV 等;采用 Catalyst 优化器进行查询优化,同时通过 Tungsten 引擎提升查询执行效率;支持多种编程语言;能够统一访问结构化和非结构化数据
dataframe 是类似于关系数据库表的分布式数据集合,提供了功能强大的数据操作方式,dataset 是在 dataframe 的基础上引入更强的 类型化的 api,可以在编译时进行类型检查,提供更好的错误检测机制和优化空间;
catalyst 是 spark sql 的查询优化器,采用规则推理和模式匹配的方式进行查询优化,主要包括逻辑计划、物理计划、优化规则等;
Tungsten 引擎旨在通过紧密控制内存和 cpu,提高 spark sql 的性能,包括高效的内存管理、无 gc 的执行模型和基于矢量化的处理模式;
spark sql 与 apache hive 具有很好的兼容性,可以直接运行 hive 的查询,并对 hive 元数据进行集成和访问
2、在 Spark SQL 中,如何创建 DataFrame?DataFrame 与 RDD 有什么区别?
a、从已有的 rdd 转换到 dataframe,需要提前定义好 schema
b、从本地的 csv、json、parquet 等文件格式读取数据
c、从数据库读取,可以通过 jdbc 连接数据库,读取表内容创建 dataframe
rdd 和 dataframe 的区别在于:
a、dataframe 是带有 schema 的分布式数据集,而 rdd 是分布式数据集合,没有结构信息
b、dataframe 有 catalyst 优化器,可以进行查询优化,而 rdd 没有这种优化机制
c、dataframe 提供了基于 sql 的高层次 api,可进行复杂的数据处理,而 rdd 提供给的是低级别的函数式编程 api
# 从 rdd 转换为 dataframe:
import spark.implicits._
val rdd = spark.sparkContext.parallelize(Seq((1, "Alice"), (2, "Bob")))
val df = rdd.toDF("id", "name")
# 从本地文件加载数据:
val dfCSV = spark.read.option("header", "true").csv("path/to/file.csv")
val dfJSON = spark.read.json("path/to/file.json")
val dfParquet = spark.read.parquet("path/to/file.parquet")
# 从数据库读取数据:
val url = "jdbc:mysql://localhost:3306/mydatabase"
val properties = new java.util.Properties()
properties.setProperty("user", "root")
properties.setProperty("password", "password")
val df = spark.read.jdbc(url, "mytable", properties)
3、Spark SQL 如何与 Hive 集成?如何在 Spark SQL 中查询 Hive 表?
a、配置 hive 环境:确保 hive 已经在执行的节点上正确安装和配置好了
b、配置 spark 环境
c、配置 hive-site.xml:将 hive 的 hive-site.xml 文件拷贝到 spark 的配置目录,确保 spark 能访问 hive 的配置
d、启动 SparkSession:传入 hive 支持参数
e、查询 hive 表:通过 spark sql 中的 sql 语句来查询 hive 表
import org.apache.spark.sql.SparkSession
val spark = SparkSession.builder
.appName("Spark Hive Example")
.config("spark.master", "local")
.config("spark.sql.warehouse.dir", "hdfs://path/to/hive/warehouse")
.enableHiveSupport()
.getOrCreate()
spark.sql("SELECT * FROM hive_table_name").show()
4、在 Spark SQL 中,如何使用 SQL 查询 DataFrame?
通过创建一个临时视图将 dataframe 注册,然后使用标准的 sql 查询这个 dataframe:
# a、导入 spark sql 库
from pyspark.sql import SparkSession
# b、创建 SparkSession 对象
spark = SparkSession.builder.appName("exampleApp").getOrCreate()
# c、读入数据并创建 DataFrame
df = spark.read.csv("data.csv", header=True, inferSchema=True)
# d、将 DataFrame 注册为临时视图
df.createOrReplaceTempView("tempView")
# e、使用 SQL 查询临时视图
result_df = spark.sql("SELECT * FROM tempView WHERE some_column > 100")
result_df.show()
5、在 Spark SQL 中,如何定义和注册一个临时视图(Temporary View)?
spark sql 中可以使用 createOrReplaceTempView 方法来定义和注册一个临时视图:首先,确保已经将数据载入到一个 Dataframe 中,然后调用 dataframe 的 createOrReplaceTempView 方法,并为该视图指定一个名称
// Step 1: 创建SparkSession
val spark = SparkSession.builder()
.appName("Temporary View Example")
.getOrCreate()
// Step 2: 载入数据到DataFrame
val df = spark.read.json("examples/src/main/resources/people.json")
// Step 3: 定义和注册临时视图
df.createOrReplaceTempView("people")
// Step 4: 使用SQL查询临时视图
val sqlDF = spark.sql("SELECT * FROM people")
// 显示结果
sqlDF.show()
在 spark sql 中,临时视图仅在创建视图的会话内有效,而全局临时视图则在所有会话中可见,直到所有会话结束
// Step 1: 创建全局临时视图
df.createGlobalTempView("people")
// Step 2: 使用全局临时视图 (需要使用 global_temp 数据库前缀)
val sqlDF = spark.sql("SELECT * FROM global_temp.people")
// 显示结果
sqlDF.show()
6、Spark SQL 中的 Catalyst 优化器是什么?它的作用是什么?
catalyst 是 spark sql 的查询优化引擎,它在查询执行之前对 sql 查询进行优化,包括逻辑计划优化、物理计划生成以及优化
a、逻辑计划优化:catalyst 首先将 sql 查询语句解析为一个初步的逻辑执行计划,随后它会进行一系列的逻辑优化,比如谓词下推、子查询去除等等,来减少数据处理量,更高效的查询
b、物理计划生成与优化:在优化过的逻辑计划的基础上,catalyst 会生成多个候选的物理执行计划,它会选择一个代价最低的执行计划,代价评估基于统计信息,如数据分布、列的基数等等
c、选定最优的物理执行计划后,catalyst 会对操作进行代码生成,执行使用 scala 代码来执行物理计划,直接操作 jvm 字节码
d、采用规则驱动机制来进行优化,规则是由开发定义的一系列模式匹配和重写规则,用于查询计划的转换
7、在 Spark SQL 中,如何通过 UDF(用户自定义函数)扩展 SQL 功能?
udf 是用户定义的,将一个或多个输入列映射到一个输出列的函数,创建 udf 的大致步骤如下:
a、定义一个函数并使用 functions.udf 方法将其转换成 udf
b、将这个 udf 注册到 sparksession,然后通过 sql 来调用它
c、在 dataframe api 中使用 withcolumn、select 等方法来调用这个 udf
# 举例:有一个字符串逆序的函数:
def reverse_string(s: str) -> str:
return s[::-1]
# 将这个函数转换为 udf:
from pyspark.sql import functions as F
reverse_udf = F.udf(reverse_string, StringType())
# 将 udf 注册到 sparksession:
spark.udf.register("reverse", reverse_udf)
# 在 sql 查询中使用这个 udf:
df.createOrReplaceTempView("my_table")
result = spark.sql("SELECT reverse(name) as reversed_name FROM my_table")
result.show()
# 在 dataframe api 中使用 udf:
df.withColumn("reversed_name", reverse_udf(df["name"])).show()
在 spark sql 中使用 udf 有一些注意事项:
a、udf 的执行速度可能比内建函数慢,因为不能很好的利用 spark 计划优化器和执行引擎的优化
b、除了 udf,用户还可以自定义聚合函数 udaf 用于与聚集操作,比如计算分组数据的标准差
c、spark 3.0 引入了 pandas udf,可以显著提高性能,更好的进行矢量化计算,直接使用 pandas 的函数来处理 spark 数据中的块数据。
from pyspark.sql.functions import pandas_udf
@pandas_udf(StringType())
def reverse_udf(s: pd.Series) -> pd.Series:
return s.apply(lambda x: x[::-1])
df.withColumn("reversed_name", reverse_udf("name")).show()
8、如何在 Spark SQL 中进行数据的分区操作?分区对性能的影响是什么?
分区操作可以优化数据的读取和处理性能,要进行分区操作,可以通过以下的方式:
a、使用 dataframwriter 的 partitionBy 方法,在写入数据时定义分区列
b、手动指定分区:使用 repartition 或 coalesce 方法来调整 dataframe 的分区数量
c、通过 sql 创建表时,使用 partitioned by 语句定义分区
分区对性能的影响主要体现在以下的几个方面:
a、数据倾斜
b、增加分区数量可以提高并行度,从而提高任务的执行速度,但是调度开销也可能变大
c、I/O 性能:分区后的数据在查询时可以减小扫描的数据量,尤其是在分区列上进行筛选查询时显著提高查询性能
分桶:进一步对数据进行分组存储,不同于分区的是分桶可以在读取时减少数据的扫描量,但是维护的开销更大,适用于需要进行高效 join 或 groupBy 操作的场景 -- 可以使用 bucketBy 方法来进行分桶
分区裁剪:在运行查询时,spark sql 会自动根据查询条件裁剪掉不必要扫描的分区
动态分区插入:在插入数据时动态生成分区,对于处理流式数据或没有预定义分区值的数据非常有用
9、Spark SQL 中的 Schema 是如何定义的?如何动态推断 Schema?
schema 是用于定义 dataframe 或 dataset 中数据的表的结构,可以通过编程方式手动定义 schema,也可以通过动态推断来自动生成 schema
a、手动定义 schema 时,通常会用到一下两个类:StructType 表示结构类型,即一组字段;StructField 表示字段,即一个字段名、字段类型和是否可以为 null 的三元组
import org.apache.spark.sql.types.{StructType, StructField, StringType, IntegerType}
val schema = StructType(List(
StructField("name", StringType, true),
StructField("age", IntegerType, true)
))
# 或使用手动定义的 schema 读取 csv 文件:
val df = spark.read
.format("csv")
.schema(schema)
.load("path/to/csvfile")
# 动态推测 schema:
val df = spark.read
.json("path/to/jsonfile")
val df = spark.read
.parquet("path/to/parquetfile")
val df = spark.read
.option("header", "true")
.csv("path/to/csvfile")
10、在 Spark SQL 中,如何使用 DataFrame API 实现复杂的查询和聚合操作?
a、创建 SparkSession
b、加载数据并将其转化为 DataFrame
c、使用 DataFrame 提供的多种操作符进行查询、过滤、聚合等复杂操作
d、最后,通过特定的动作操作将结果输出
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
# 创建 SparkSession
spark = SparkSession.builder.appName("ComplexQueryExample").getOrCreate()
# 读取数据,假设数据存储为 JSON 格式
df = spark.read.json("path/to/your/data.json")
# 进行复杂查询和聚合操作
result = df.filter(col('age') > 21) \
.groupBy('gender') \
.agg(avg('income').alias('average_income'), max('income').alias('max_income')) \
.orderBy('average_income', ascending=False)
# 显示结果
result.show()
dataframe api 还提供了丰富的功能,比如 连接操作:内连接、外连接、左连接、右连接等
df1 = spark.read.json("path/to/data1.json")
df2 = spark.read.json("path/to/data2.json")
joined_df = df1.join(df2, df1["id"] == df2["id"], "inner")
窗口函数:排名、移动、平均等:
from pyspark.sql.window import Window
windowSpec = Window.partitionBy('gender').orderBy('income')
df = df.withColumn('rank', rank().over(windowSpec))
df.show()
自定义 udf:
from pyspark.sql.types import IntegerType
def increment_age(age):
return age + 1
spark.udf.register("incrementAge", increment_age, IntegerType())
df = df.withColumn('new_age', udf(lambda x: increment_age(x), IntegerType())(df['age']))
df.show()
缓存和持久化:为了提高查询性能,可以对 DataFrame 进行缓存或持久化
df.cache() # 或者 df.persist()
result = df.filter(col('age') > 21).groupBy('gender').agg(avg('income').alias('average_income'))
result.show()
11、Spark SQL 中的 DataSet 和 DataFrame 有什么区别?如何选择使用?
dataset 和 dataframe 分别是处理结构化数据的两种主要抽象,有一些核心的区别:
a、数据类型安全:dataframe 是分布式数据集的抽象,相当于一个分布式的二维表,其列是用命名表示的,更接近 sql 表或者 pandas dataframe,只能以字符串形式来🚰列名,这样就导致了类型安全丧失;
dataset 类似于 java/scala 里的强类型集合,可以用类型化的操作如 map、filter等来处理,是结合了 rdd 的强类型和 dataframe 的优化优势;
b、dataframe 注重对数据的 sql 查询,主要提供了 select、filter 等高阶函数;dataset 提供了数据库风格的操作,也可以提供类型安全的操作,如使用 lambda 表达式进行过滤等
c、dataframe 经过 catalyst 优化器进行优化后,后生成更加高效的查询计划,执行速度较快,并且进行 Tungsten 物理执行计划优化;dataset 也经过 catalyst 优化器优化,但由于其强类型的特性,某些场景下它的性能可能比 dataframe 要好;
d、如果需要简洁的 sql 风格查询操作,可以选择 dataframe,如果需要类型安全的、强类型的操作,选择 dataset;dataframe api 常用于查询分析,在一些应用中效率可能更高,dataset 在对大规模和复杂数据处理中优势明显,因为类型安全校验降低了错误概率
12、Spark SQL 是如何优化查询计划的?Explain 语句的作用是什么?
a、解析:将 sql 文本解析成 Unresolved Logical Plan
b、分析:通过解析表、列等元数据信息,生成 Resolved Logical Plan
c、优化:对已解决的逻辑计划应用一系列规则进行优化,生成 Optimized Logical Plan,例如谓词下推、常量折叠、列剪裁,逻辑优化的目标是简化查询的计算,减少数据扫描量
d、物理规划:将优化后的逻辑计划转换为一系列具体的物理操作,生成物理计划,不同物理计划可能会选择不同的执行策略,比如 broadcast join 与 shuffle join
e、代码生成:生成具体执行任务的代码,例如将 rdd 转换为 spark 任务
可以通过 explain 来看具体的执行计划:
explain 输出整个优化过程的简短概述;explain extended 输出包括初始逻辑计划、优化后的逻辑计划和物理计划的详细信息;explain formatted 以更易读的格式显示查询计划,适用于需要详细分析查询计划的场景
13、Spark SQL 是如何处理内存中的大数据集的?它如何避免内存溢出?
a、spark 将内存分为 两个部分:执行内存和存储内存,这两部分动态调整,根据任务需求互相借用
b、spark 提供多种持久化级别,可以将数据持久化到内存、磁盘上或二者结合,保证了即便内存不足,可以将数据部分写入磁盘。
c、调整分区数是的数据分区更为合理,为内存管理和数据处理提供良好负载均衡
d、kryo 序列化提高内存存储和传输的效率
e、查询优化 catalyst
14、Spark SQL 中的 Tungsten 优化是什么?它对性能提升的关键点是什么?
a、缓存行对齐:通过内存布局优化,将数据更紧凑地存储到内存中,以减少内存消耗和垃圾回收的成本,显著降低 cpu 缓存未命中的次数
b、将 spark sql 查询计划直接编译为 java 字节码,减少解释执行带来的开销
c、无 gc
d、对磁盘 io 和网络通信进行了优化,进一步提升了 spark 在分布式数据处理环境下的性能
15、Spark SQL 的广播连接(Broadcast Join)是什么?在什么情况下使用?
spark.sql.autoBroadcastJoinThreshold 定义了广播连接的大小阈值(默认为 10 MB),如果小表的大小小于该阈值,spark 会自动选择广播连接,也可以在 dataframe 的 join 操作中,使用 broadcast 方法来强制广播连接:val result = largeDF.join(org.apache.spark.sql.functions.broadcast(smallDF), "key")
16、在 Spark SQL 中,如何实现窗口函数操作?常见的窗口函数有哪些?
通过 window 对象和 partitionBy,orderBy 等方法,可以定义不同的窗口范围;
常见窗口函数:
rank:根据排序表达式为每行分配等级,有相同值跳过中间值
dense_rank:类似 rank,但不会跳过中间值
row_number:为每行分配一个唯一的行编号
lead:获取当前行之后指定第几行的值
lag:获取当前行之前指定第几行的值
cume_dist:计算相对累计分布
ntile:将数据按照某个等分数划分
使用场景:排名、排序、移动平均和累积和、前后数据差异
17、Spark SQL 如何通过缓存(Cache)提高查询效率?缓存机制的作用是什么?
使用 persist() 和 cache() 方法来实现,cache 是 memory_only,persist 可以指定数据存储级别,缓存机制的作用是提升查询速度、降低 I/O 开销、提高资源利用率
什么时候使用缓存:a、需要重复多次使用的数据 b、长时间重算的数据 c、处理速度对用户体验影响较大的场景
18、在 Spark SQL 中,如何处理数据倾斜问题?有哪些优化策略?
调整并行度、使用 salting、自定义分区器、优化数据读取(通过 分区表、bucket 表)、广播小表、skew join(针对一些特别热的 key,采用单独处理的方式)
19、如何在 Spark SQL 中进行表的分区和分桶?两者的区别是什么?
分区是按列的值将数据拆分成多个较小的部分,每个分区都会存放在单独的文件或目录下,查询时只需扫描相关分区,减少数据读取量
CREATE TABLE student (
name STRING,
age INT,
grade STRING
)
PARTITIONED BY (grade STRING);
分桶:将表中的数据按照某个列的 hash 值进行分组,然后把这些数据分到固定数量的桶里,每个桶会包含一个相对均匀的数据子集
CREATE TABLE student_buckets (
name STRING,
age INT,
grade STRING
)
CLUSTERED BY (age) INTO 4 BUCKETS;
分区主要是为了减少扫描时的数据读取量,适用于按值分开的数据;分桶时为了数据的均匀分布与高效的采样,且桶的数量固定
分区可以较少 io,因为只需扫描必要的分区;在 join 操作中,桶表之间的 join 操作可以并行执行,提升效率;分区和分桶可以一起使用,比如按年分区,再按月分桶
20、Spark SQL 中的 SQL 查询与 DataFrame API 查询有什么区别?
sql 查询使用的是 sql 语句,直观但编写复杂转换逻辑不灵活;dataframe api 使用的是 spark 提供的高级编程接口,通过链式调用进行操作,更灵活