Spark SQL 中的 SQL 查询与 DataFrame API 查询有什么区别?
·
Spark SQL 中的 SQL 查询和 DataFrame API 查询在功能上是等价的,但在使用方式、表达能力和优化机制上有重要区别。
1. 基本语法对比
1.1 SQL 查询方式
-- 注册临时视图
CREATE OR REPLACE TEMPORARY VIEW sales_2025_sep AS
SELECT * FROM parquet.`/data/sales_2025_sep.parquet`;
-- 标准SQL查询:分析2025年9月销售数据
SELECT
product_category,
DATE_FORMAT(sale_date, 'yyyy-MM-dd') as sale_day,
COUNT(*) as daily_transactions,
AVG(amount) as avg_amount,
SUM(amount) as daily_sales
FROM sales_2025_sep
WHERE sale_date BETWEEN '2025-09-01' AND '2025-09-30'
AND region IN ('north', 'south')
GROUP BY product_category, DATE_FORMAT(sale_date, 'yyyy-MM-dd')
HAVING daily_sales > 5000
ORDER BY sale_day, daily_sales DESC;
1.2 DataFrame API 查询方式
import org.apache.spark.sql.functions._
import java.time.LocalDate
// 读取2025年9月数据
val salesDF = spark.read.parquet("/data/sales_2025_sep.parquet")
// 定义2025年9月的时间范围
val sepStart = LocalDate.of(2025, 9, 1)
val sepEnd = LocalDate.of(2025, 9, 30)
// DataFrame API 链式调用
val resultDF = salesDF
.filter(col("sale_date").between(lit(sepStart.toString), lit(sepEnd.toString)))
.filter(col("region").isin("north", "south"))
.withColumn("sale_day", date_format(col("sale_date"), "yyyy-MM-dd"))
.groupBy("product_category", "sale_day")
.agg(
count("*").as("daily_transactions"),
avg("amount").as("avg_amount"),
sum("amount").as("daily_sales")
)
.filter(col("daily_sales") > 5000)
.orderBy("sale_day", desc("daily_sales"))
2. 核心区别对比
2.1 语法和表达方式
| 特性 | SQL 查询 | DataFrame API |
|---|---|---|
| 语法风格 | 声明式,类似传统SQL | 链式方法调用,面向对象 |
| 类型安全 | 运行时类型检查 | 编译时类型检查(Scala) |
| 代码补全 | 有限支持 | 完整的IDE支持 |
| 重构能力 | 较弱 | 强大的重构支持 |
2.2 类型安全示例
// DataFrame API - 编译时类型检查
val salesDF: DataFrame = spark.read.parquet("/data/sales_2025_sep.parquet")
// 编译错误:列名拼写错误会被IDE检测
// salesDF.select("produc_category") // 编译时报错
// 正确写法 - 编译时安全
salesDF.select(col("product_category"))
// SQL方式 - 运行时才报错
spark.sql("SELECT produc_category FROM sales_2025_sep") // 运行时异常
3. 复杂查询场景对比
3.1 窗口函数应用(2025年9月数据分析)
-- SQL方式:计算2025年9月每日销售额的7日移动平均
SELECT
sale_date,
product_category,
daily_sales,
AVG(daily_sales) OVER (
PARTITION BY product_category
ORDER BY sale_date
ROWS BETWEEN 6 PRECEDING AND CURRENT ROW
) as moving_avg_7d
FROM (
SELECT
sale_date,
product_category,
SUM(amount) as daily_sales
FROM sales_2025_sep
WHERE sale_date BETWEEN '2025-09-01' AND '2025-09-30'
GROUP BY sale_date, product_category
) daily_stats
ORDER BY product_category, sale_date;
// DataFrame API方式:相同的窗口计算
import org.apache.spark.sql.expressions.Window
val windowSpec = Window
.partitionBy("product_category")
.orderBy("sale_date")
.rowsBetween(-6, 0)
val dailyStats = salesDF
.filter(col("sale_date").between("2025-09-01", "2025-09-30"))
.groupBy("sale_date", "product_category")
.agg(sum("amount").as("daily_sales"))
val resultWithMovingAvg = dailyStats
.withColumn("moving_avg_7d",
avg("daily_sales").over(windowSpec))
.orderBy("product_category", "sale_date")
3.2 多表连接查询
-- SQL方式:连接2025年9月销售数据和产品维度表
SELECT
s.sale_date,
p.product_name,
p.category,
s.amount,
c.customer_segment
FROM sales_2025_sep s
JOIN products p ON s.product_id = p.product_id
JOIN customers c ON s.customer_id = c.customer_id
WHERE s.sale_date BETWEEN '2025-09-01' AND '2025-09-30'
AND p.is_active = true
AND c.registration_date <= '2025-08-31';
// DataFrame API方式:相同的多表连接
val productsDF = spark.table("products")
val customersDF = spark.table("customers")
val joinedResult = salesDF
.filter(col("sale_date").between("2025-09-01", "2025-09-30"))
.join(productsDF.filter(col("is_active") === true), "product_id")
.join(customersDF.filter(col("registration_date") <= "2025-08-31"), "customer_id")
.select(
col("sale_date"),
col("product_name"),
col("category"),
col("amount"),
col("customer_segment")
)
4. 执行计划和优化对比
4.1 执行计划相同性
// 两种方式最终生成相同的执行计划
val sqlQuery = """
SELECT product_category, SUM(amount) as total
FROM sales_2025_sep
WHERE sale_date BETWEEN '2025-09-01' AND '2025-09-30'
GROUP BY product_category
"""
val dfQuery = salesDF
.filter(col("sale_date").between("2025-09-01", "2025-09-30"))
.groupBy("product_category")
.agg(sum("amount").as("total"))
// 比较执行计划 - 两者完全相同
println("SQL执行计划:")
spark.sql(sqlQuery).explain()
println("\nDataFrame执行计划:")
dfQuery.explain()
4.2 Catalyst优化器工作流程
5. 开发体验对比
5.1 代码可读性和维护性
// DataFrame API - 更好的模块化和复用
def createSeptember2025Filter(): Column = {
col("sale_date").between("2025-09-01", "2025-09-30")
}
def calculateSalesMetrics(): Seq[Column] = Seq(
count("*").as("transaction_count"),
avg("amount").as("avg_amount"),
sum("amount").as("total_sales")
)
// 可复用的查询构建
val baseQuery = salesDF.filter(createSeptember2025Filter())
val categoryAnalysis = baseQuery
.groupBy("product_category")
.agg(calculateSalesMetrics(): _*)
val regionalAnalysis = baseQuery
.groupBy("region")
.agg(calculateSalesMetrics(): _*)
5.2 动态查询构建
// DataFrame API - 动态构建复杂查询
def buildDynamicQuery(filters: Map[String, Any]): DataFrame = {
var query = salesDF.filter(col("sale_date").between("2025-09-01", "2025-09-30"))
filters.foreach { case (column, value) =>
value match {
case list: Seq[_] => query = query.filter(col(column).isin(list: _*))
case singleValue => query = query.filter(col(column) === singleValue)
}
}
query
}
// 动态应用过滤器
val dynamicResult = buildDynamicQuery(Map(
"region" -> Seq("north", "south"),
"payment_method" -> "credit_card"
)).groupBy("product_category").agg(sum("amount").as("total"))
6. 性能考虑
6.1 缓存策略应用
// DataFrame API可以更精细地控制缓存
val septemberData = salesDF
.filter(col("sale_date").between("2025-09-01", "2025-09-30"))
.cache() // 精确缓存需要的数据
// 多个查询共享缓存
val result1 = septemberData.groupBy("product_category").agg(sum("amount"))
val result2 = septemberData.groupBy("region").agg(avg("amount"))
// SQL方式缓存整个表
spark.sql("CACHE TABLE sales_2025_sep")
6.2 分区剪裁优化
// 如果数据按日期分区,两者都能利用分区剪裁
// 分区目录结构: /sale_date=2025-09-01/, /sale_date=2025-09-02/, etc.
// DataFrame API
salesDF.filter(col("sale_date").between("2025-09-01", "2025-09-15"))
// SQL方式
spark.sql("SELECT * FROM sales_2025_sep WHERE sale_date BETWEEN '2025-09-01' AND '2025-09-15'")
7. 实际应用建议
7.1 选择标准
def chooseApproach(scenario: String): String = scenario match {
case "简单查询" | "临时分析" => "SQL"
case "复杂业务逻辑" | "需要类型安全" => "DataFrame API"
case "ETL管道" | "生产代码" => "DataFrame API"
case "即席查询" | "数据探索" => "SQL"
case _ => "根据团队习惯选择"
}
7.2 混合使用最佳实践
// 结合两者优势
// 1. 使用DataFrame API进行数据准备
val preparedData = salesDF
.filter(col("sale_date").between("2025-09-01", "2025-09-30"))
.select("product_id", "amount", "customer_id")
.cache()
// 2. 注册为临时视图进行复杂SQL分析
preparedData.createOrReplaceTempView("sep_2025_prepared")
// 3. 使用SQL进行复杂分析
val complexAnalysis = spark.sql("""
WITH customer_stats AS (
SELECT
customer_id,
COUNT(*) as order_count,
SUM(amount) as total_spent
FROM sep_2025_prepared
GROUP BY customer_id
)
SELECT
CASE
WHEN total_spent > 1000 THEN 'VIP'
WHEN total_spent > 500 THEN 'Regular'
ELSE 'New'
END as segment,
AVG(order_count) as avg_orders
FROM customer_stats
GROUP BY
CASE
WHEN total_spent > 1000 THEN 'VIP'
WHEN total_spent > 500 THEN 'Regular'
ELSE 'New'
END
""")
总结
SQL查询的优势:
- 语法熟悉,学习成本低
- 适合即席查询和数据分析
- 与现有SQL工具兼容性好
DataFrame API的优势:
- 编译时类型安全
- 更好的IDE支持
- 代码可读性和可维护性更强
- 更适合复杂业务逻辑和ETL管道
更多推荐


所有评论(0)