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优化器工作流程

SQL字符串
SQL解析
DataFrame操作
逻辑计划生成
Unresolved Logical Plan
分析Analysis
逻辑优化
物理计划
代码生成
执行

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管道
Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐