Spark SQL 的 CBO(Cost-Based Optimizer)通过统计信息来估算不同执行计划的代价,从而选择最优的执行策略。以下是详细的优化方法和配置:

1. CBO 基础配置与启用

核心配置参数

// 启用 CBO 核心功能
spark.conf.set("spark.sql.cbo.enabled", "true")
spark.conf.set("spark.sql.cbo.joinReorder.enabled", "true")

// 统计信息相关配置
spark.conf.set("spark.sql.statistics.histogram.enabled", "true")
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "50MB")

// 自适应查询执行(AQE)与 CBO 协同工作
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")

2. 统计信息收集与管理

表级统计信息收集

-- 收集全表统计信息
ANALYZE TABLE sales COMPUTE STATISTICS;

-- 收集特定列的统计信息
ANALYZE TABLE sales COMPUTE STATISTICS FOR COLUMNS 
    product_id, amount, sale_date;

-- 收集分区表统计信息
ANALYZE TABLE sales_partitioned PARTITION(dt='2023-01-01') COMPUTE STATISTICS;

通过 DataFrame API 收集统计信息

// 创建表时收集统计信息
val df = spark.read.parquet("data/sales.parquet")
df.write.option("computeStats", "true").saveAsTable("sales")

// 手动触发统计信息收集
spark.sql("ANALYZE TABLE sales COMPUTE STATISTICS")

// 验证统计信息
spark.sql("DESCRIBE EXTENDED sales").show(truncate=false)

统计信息内容解析

// 查看表的详细统计信息
val stats = spark.sql("DESCRIBE EXTENDED sales")
stats.filter(col("col_name").contains("Statistics")).show(false)

/* 输出示例:
+------------------+------------------------------------------------+
|col_name          |data_type                                       |
+------------------+------------------------------------------------+
|Statistics        |123456 bytes, 1000000 rows                      |
+------------------+------------------------------------------------+
*/

3. CBO 优化场景实战

Join 顺序优化

-- 原始查询(可能不是最优顺序)
EXPLAIN COST
SELECT * 
FROM large_table l
JOIN medium_table m ON l.id = m.large_id
JOIN small_table s ON m.small_id = s.id;

-- CBO 会自动重新排序 Join 顺序
/* 优化后的计划可能:
BroadcastHashJoin [small_id], [id], Inner, BuildRight  -- 先广播小表
:- BroadcastHashJoin [id], [large_id], Inner, BuildRight
:  :- TableScan large_table
:  +- BroadcastExchange
:     +- TableScan medium_table
+- BroadcastExchange
   +- TableScan small_table
*/

Join 类型选择优化

// 不同数据量下的 Join 策略选择
val largeDF = spark.table("large_table")  // 1亿行
val mediumDF = spark.table("medium_table") // 100万行  
val smallDF = spark.table("small_table")   // 1千行

// CBO 会根据统计信息自动选择:
// - smallDF: BroadcastHashJoin
// - mediumDF: SortMergeJoin 或 ShuffleHashJoin
// - largeDF: SortMergeJoin

val result = largeDF.join(mediumDF, "id").join(smallDF, "id")
result.explain("cost")

4. 复杂查询优化案例

多表关联优化

-- 复杂业务查询示例
WITH customer_stats AS (
    SELECT 
        customer_id,
        COUNT(*) as order_count,
        SUM(amount) as total_amount
    FROM orders 
    WHERE order_date >= '2023-01-01'
    GROUP BY customer_id
),
product_stats AS (
    SELECT 
        product_id,
        AVG(price) as avg_price,
        COUNT(DISTINCT category) as category_count
    FROM products 
    GROUP BY product_id
)
SELECT 
    c.customer_name,
    cs.order_count,
    cs.total_amount,
    p.product_name,
    ps.avg_price
FROM customers c
JOIN customer_stats cs ON c.customer_id = cs.customer_id
JOIN orders o ON c.customer_id = o.customer_id  
JOIN products p ON o.product_id = p.product_id
JOIN product_stats ps ON p.product_id = ps.product_id
WHERE cs.total_amount > 1000
  AND ps.avg_price > 50;

-- CBO 优化效果:
-- 1. 自动选择最优的 Join 顺序
-- 2. 根据数据量选择合适的 Join 算法
-- 3. 谓词下推减少中间结果集

子查询优化

-- 相关子查询优化
SELECT 
    employee_id,
    salary,
    (SELECT AVG(salary) FROM employees e2 
     WHERE e2.department = e1.department) as dept_avg_salary
FROM employees e1
WHERE salary > 50000;

-- CBO 可能重写为更高效的 Join 操作
/* 优化后等价于:
SELECT 
    e1.employee_id,
    e1.salary,
    dept_avg.avg_salary as dept_avg_salary
FROM employees e1
JOIN (
    SELECT department, AVG(salary) as avg_salary 
    FROM employees 
    GROUP BY department
) dept_avg ON e1.department = dept_avg.department
WHERE e1.salary > 50000;
*/

5. 高级优化技巧

直方图统计优化

// 启用直方图统计(对数据分布不均匀的列特别有效)
spark.conf.set("spark.sql.statistics.histogram.enabled", "true")
spark.conf.set("spark.sql.statistics.histogram.numBins", "254")

// 为关键列收集直方图统计信息
spark.sql("ANALYZE TABLE sales COMPUTE STATISTICS FOR COLUMNS amount, product_id")

// 直方图帮助 CBO 更好地估算:
// - 数据倾斜程度
// - 选择性估计(WHERE 条件过滤效果)
// - Join 结果集大小预估

谓词选择性估计

-- CBO 利用统计信息估计谓词的选择性
SELECT * FROM sales 
WHERE amount > 1000 
  AND product_id IN ('A001', 'A002', 'A003')
  AND sale_date BETWEEN '2023-01-01' AND '2023-03-31';

-- CBO 会基于以下信息优化:
-- 1. amount 列的 min/max/ndv(不同值数量)
-- 2. product_id 的直方图分布
-- 3. sale_date 的范围统计

6. 监控与诊断

执行计划分析

val df = spark.sql("""
SELECT c.name, SUM(o.amount) as total
FROM customers c
JOIN orders o ON c.customer_id = o.customer_id  
JOIN products p ON o.product_id = p.product_id
WHERE o.order_date >= '2023-01-01'
  AND p.category = 'ELECTRONICS'
GROUP BY c.name
""")

// 查看详细的代价优化计划
df.explain("cost")
df.explain("formatted")

// 重点关注:
// - 是否出现 Cost 相关的优化器决策
// - Join 顺序是否合理
// - 是否存在 BroadcastExchange

CBO 效果验证

// 对比开启/关闭 CBO 的性能差异
def benchmarkQuery(spark: SparkSession, query: String): Long = {
  val startTime = System.currentTimeMillis()
  spark.sql(query).count()
  System.currentTimeMillis() - startTime
}

// 测试 CBO 效果
spark.conf.set("spark.sql.cbo.enabled", "false")
val timeWithoutCBO = benchmarkQuery(spark, complexQuery)

spark.conf.set("spark.sql.cbo.enabled", "true")  
val timeWithCBO = benchmarkQuery(spark, complexQuery)

println(s"Without CBO: ${timeWithoutCBO}ms, With CBO: ${timeWithCBO}ms")

7. 最佳实践总结

CBO 优化检查清单

def setupCBOOptimizations(spark: SparkSession): Unit = {
  // 基础 CBO 配置
  spark.conf.set("spark.sql.cbo.enabled", "true")
  spark.conf.set("spark.sql.cbo.joinReorder.enabled", "true")
  
  // 统计信息配置
  spark.conf.set("spark.sql.statistics.histogram.enabled", "true") 
  spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "50MB")
  
  // AQE 协同优化
  spark.conf.set("spark.sql.adaptive.enabled", "true")
  spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
  
  // 定期收集统计信息
  spark.sql("ANALYZE TABLE sales COMPUTE STATISTICS")
}

// 表维护脚本示例
def maintainTableStats(spark: SparkSession, tableName: String): Unit = {
  spark.sql(s"ANALYZE TABLE $tableName COMPUTE STATISTICS")
  spark.sql(s"ANALYZE TABLE $tableName COMPUTE STATISTICS FOR COLUMNS " + 
    "id, amount, date_col, category")
}

数据特征与优化策略匹配

小表<阈值
中等表
大表
数据特征分析
数据分布均匀?
标准CBO优化
直方图统计优化
表大小对比?
广播Join优化
SortMergeJoin优化
分区裁剪优化
收集统计信息验证

通过系统性地应用 CBO 优化策略,Spark SQL 能够智能地选择最优的执行计划,特别是在复杂多表关联和数据分布不均匀的场景下,性能提升效果显著。

Logo

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

更多推荐