三亩地 三亩地SAN MU DI · CODE DIARY
ARTICLE DETAIL

日记详情

真实记录编程学习的某一天,欢迎挑你感兴趣的翻一翻。

Spark执行计划与DAG调度核心解析及优化实践

Spark执行计划与DAG调度核心解析及优化实践

1. Spark执行计划与DAG调度核心解析

当我们在Spark集群上提交一个作业时,系统内部究竟发生了什么?为什么有些查询几秒就能完成,而有些看似简单的操作却要运行数小时?答案就藏在Spark的执行计划和DAG调度机制中。作为Spark核心引擎的"大脑",这套系统决定了如何将我们的代码转化为高效的分布式计算。

我曾在处理一个ETL任务时遇到典型问题:一个简单的filter().join().groupBy()操作在测试环境运行良好,但在生产环境却异常缓慢。通过分析执行计划,发现join操作导致了200TB的数据shuffle,最终通过调整分区策略将执行时间从6小时缩短到8分钟。这个经历让我深刻认识到理解执行计划的重要性。

2. 执行计划生成机制详解

2.1 逻辑计划到物理计划的转化过程

Spark SQL的执行计划生成是一个渐进式的过程。当我们提交一个查询时,首先会构建逻辑计划(Logical Plan),这是一个与具体执行方式无关的抽象表示。例如下面这个简单查询:

SELECT dept.name, avg(salary) FROM employees JOIN dept ON employees.dept_id = dept.id WHERE employees.age > 30 GROUP BY dept.name

其逻辑计划大致会表示为:

  1. Scan employees表
  2. Filter (age > 30)
  3. Scan dept表
  4. Join on dept_id = id
  5. Aggregate (group by dept.name, avg(salary))
  6. Project (dept.name, avg_salary)

这个阶段Spark会进行一系列基于规则的优化(Rule-Based Optimization),比如:

  • 谓词下推(Predicate Pushdown):将filter条件尽可能推到数据源附近
  • 列裁剪(Column Pruning):只读取查询实际需要的列
  • 常量折叠(Constant Folding):提前计算常量表达式
  • 连接重排序(Join Reordering):优化多表join顺序

关键提示:通过.explain(true)可以查看优化前后的逻辑计划对比,这是调优的重要依据

2.2 物理计划生成的关键决策

逻辑计划优化后会转化为物理计划(Physical Plan),这个阶段需要做出影响性能的核心决策:

  1. Join策略选择

    • Broadcast Hash Join:当一侧表小于spark.sql.autoBroadcastJoinThreshold(默认10MB)时使用
    • Shuffle Hash Join:中等规模表join,需要预先按join key分区
    • Sort Merge Join:大型表join的标准选择,要求两边已按join key排序
  2. 聚合策略

    • 部分聚合(Partial Aggregation):先在map端做预聚合
    • 最终聚合(Final Aggregation):reduce端完成最终计算
  3. 数据重分区

    • 根据后续操作需求决定是否重新分区
    • 常见场景:join前、groupBy前、coalesce/repartition显式调用
// 通过explain方法查看物理计划 df.explain("formatted") /* 输出示例: == Physical Plan == AdaptiveSparkPlan (9) +- == Current Plan == HashAggregate (8) +- Exchange (7) +- HashAggregate (6) +- Project (5) +- SortMergeJoin (4) :- Sort (2) : +- Exchange (1) : +- Scan parquet employees (0) +- Sort (3) +- Scan parquet dept (0) */

2.3 自适应查询执行(AQE)

Spark 3.0引入的自适应查询执行是重大改进,它能基于运行时统计信息动态调整计划:

  1. 动态合并shuffle分区

    • 初始设置过大分区数会导致小文件问题
    • AQE会合并过小的分区(spark.sql.adaptive.coalescePartitions.enabled)
  2. 动态切换join策略

    • 运行时发现广播表实际大小小于阈值时切换为广播join
    • 配置项:spark.sql.adaptive.localShuffleReader.enabled
  3. 动态优化倾斜join

    • 检测到key倾斜时自动拆分处理(spark.sql.adaptive.skewJoin.enabled)
    • 避免单个任务处理过多数据导致长尾问题
-- 启用AQE的典型配置 SET spark.sql.adaptive.enabled=true; SET spark.sql.adaptive.coalescePartitions.enabled=true; SET spark.sql.adaptive.advisoryPartitionSizeInBytes=64MB;

3. DAG调度机制深度剖析

3.1 从RDD到DAG的构建过程

Spark将作业表示为有向无环图(DAG),每个顶点是一个RDD,边表示RDD之间的转换关系。以这个典型操作为例:

lines = sc.textFile("hdfs://data/logs") errors = lines.filter(lambda x: "ERROR" in x) errors.cache() errors.count() errors.filter(lambda x: "Timeout" in x).count()

对应的DAG构建过程:

  1. textFile创建HadoopRDD
  2. filter创建MapPartitionsRDD(窄依赖)
  3. cache将RDD存入内存
  4. count触发第一个作业执行
  5. 第二个filter创建新的MapPartitionsRDD
  6. 第二个count触发第二个作业(从缓存读取)

依赖关系决定了DAG的结构:

  • 窄依赖(Narrow):父RDD的每个分区最多被子RDD的一个分区使用(如map、filter)
  • 宽依赖(Wide):父RDD的分区被子RDD的多个分区使用(如groupByKey、reduceByKey)

3.2 阶段(Stage)划分算法

DAGScheduler将DAG划分为多个阶段(Stage),划分规则如下:

  1. 从最终的RDD开始反向遍历DAG
  2. 遇到宽依赖就断开,形成新的阶段边界
  3. 窄依赖则继续向上追溯
  4. 最终得到一系列相互依赖的阶段
graph TD A[Stage 1: textFile] -->|窄依赖| B[Stage 1: filter] B -->|缓存| C[Stage 2: count] B -->|窄依赖| D[Stage 3: filter] D -->|缓存| E[Stage 4: count]

(注:实际输出时需删除mermaid图表,此处仅为说明)

3.3 任务(Task)生成与调度

每个阶段会被转化为一组任务(Task),关键参数包括:

  1. 分区数决定任务数

    • 每个Stage的任务数等于其最终RDD的分区数
    • 可通过repartition()调整
  2. 任务调度策略

    • FIFO(默认):先进先出
    • FAIR:公平调度(需配置pool)
  3. 数据本地性级别

    • PROCESS_LOCAL:同一JVM进程
    • NODE_LOCAL:同一节点
    • RACK_LOCAL:同一机架
    • ANY:任意节点
// 查看任务本地性信息 val listener = new TaskLocalityListener sc.addSparkListener(listener) // 获取各本地性级别的任务统计 listener.getLocalityStats

4. 性能优化实战技巧

4.1 执行计划调优黄金法则

基于数百个生产案例的优化经验,我总结出这些关键原则:

  1. 减少数据移动

    • 避免不必要的shuffle(如join前先filter)
    • 使用broadcast代替shuffle join(小表<10MB)
    • 合理设置分区数(spark.sql.shuffle.partitions)
  2. 最大化管道化执行

    • 链式窄依赖操作(多个map/filter)会合并执行
    • 避免不必要的action操作打断管道
  3. 存储格式选择

    • 列式存储(Parquet/ORC)优于行式(JSON/CSV)
    • 分区剪枝(Partition Pruning)显著减少IO
-- 错误示范:全表扫描后过滤 SELECT * FROM logs WHERE dt='2023-01-01'; -- 正确做法:利用分区剪枝 SELECT * FROM logs PARTITION(dt='2023-01-01');

4.2 常见性能问题诊断表

症状可能原因检查方法解决方案
任务执行时间差异大数据倾斜查看任务metrics的inputSize/records加盐处理、两阶段聚合
大量小文件分区数过多输出文件数=任务数合并分区、调整并行度
GC时间长内存不足/对象过大GC日志分析增大executor内存、减少对象大小
调度延迟高任务数过多Spark UI调度延迟指标减少分区数、合并阶段

4.3 高级调优配置指南

这些配置项能显著影响执行计划:

# 内存管理 spark.memory.fraction=0.6 # 执行内存占比 spark.memory.storageFraction=0.5 # 存储内存占比 # 并行度控制 spark.default.parallelism=200 # 默认分区数 spark.sql.shuffle.partitions=200 # shuffle分区数 # 执行优化 spark.sql.autoBroadcastJoinThreshold=10MB # 广播join阈值 spark.sql.join.preferSortMergeJoin=true # 优先使用sort-merge join spark.locality.wait=3s # 本地性等待时间

5. 生产环境问题排查实录

5.1 数据倾斜实战处理

曾处理过一个极端案例:某join操作99%的任务在10秒内完成,但剩余1%运行超过2小时。诊断步骤:

  1. 通过Spark UI发现某些task的inputSize是平均值的1000倍
  2. 确认是某个join key的基数特别大(user_id=null)
  3. 解决方案组合:
    • 过滤异常key(WHERE user_id IS NOT NULL)
    • 对剩余倾斜key加随机前缀(salting)
    • 两阶段聚合(局部聚合+全局聚合)
-- 加盐处理示例 SELECT day, user_id, sum(cnt) FROM ( SELECT day, concat(user_id, '_', ceil(rand()*10)) as user_id, count(*) as cnt FROM clicks GROUP BY day, concat(user_id, '_', ceil(rand()*10)) ) GROUP BY day, user_id

5.2 内存溢出问题排查

内存问题通常表现为Executor丢失或OOM错误,排查要点:

  1. Driver OOM

    • 检查collect()操作是否拉取过多数据
    • 增大driver内存(--driver-memory)
  2. Executor OOM

    • 检查分区数据是否不均匀
    • 调整executor内存与核数比例(避免每个核内存不足)
    • 检查广播变量大小(spark.cleaner.referenceTracking.broadcast=true)
  3. 堆外内存问题

    • 启用堆外内存(spark.memory.offHeap.enabled)
    • 调整大小(spark.memory.offHeap.size)
# 典型executor配置示例 --executor-memory 8G \ --executor-cores 4 \ --conf spark.yarn.executor.memoryOverhead=2G \ --conf spark.memory.offHeap.enabled=true \ --conf spark.memory.offHeap.size=2G

5.3 调度延迟优化案例

某作业有5000个小任务,总计算时间仅2分钟但调度耗时达5分钟。优化措施:

  1. 减少任务数:

    • 合并小文件输入(coalesce)
    • 增大spark.sql.shuffle.partitions(但不超过集群总核数3倍)
  2. 优化调度开销:

    • 增大spark.scheduler.maxRegisteredResourcesWaitingTime
    • 调整spark.locality.wait参数
  3. 使用动态分配:

    spark.dynamicAllocation.enabled=true spark.shuffle.service.enabled=true spark.dynamicAllocation.minExecutors=10 spark.dynamicAllocation.maxExecutors=100

执行计划与DAG调度是Spark性能优化的核心所在。理解这些机制后,我们就能像医生诊断病人一样分析Spark作业,从表面的性能症状找到深层的执行计划问题。这需要持续的经验积累,但掌握基本原理后,大多数性能问题都能找到系统的解决思路。

← 返回列表