请描述 PySpark 中任务调度的执行流程,并说明可以采取哪些措施来优化任务调度过程?
考察说明
考察对 PySpark 任务调度机制的理解以及优化调度的实践能力。
回答思路
- 【回答框架 1】PySpark 任务调度基于 DAG 执行引擎。提交作业后,SparkContext 构建 DAG,DAGScheduler 将 DAG 划分为多个 Stage,划分依据是宽依赖(ShuffleDependency)。每个 Stage 包含一组可并行执行的任务,由 TaskScheduler 将任务分发到集群的 Executor 上执行。
- 【回答框架 2】优化任务调度可从以下几方面入手:1. 减少 Shuffle:使用窄依赖操作(如 map、filter),避免不必要的宽依赖;2. 调整并行度:设置合理的分区数(如 spark.sql.shuffle.partitions),使任务数与集群资源匹配;3. 数据本地性:通过调整 spark.locality.wait 等参数,等待本地数据任务,减少网络传输;4. 资源配置:合理设置 Executor 内存、核心数,避免资源浪费或GC开销。
- 【回答框架 3】还可以通过以下方式优化:1. 使用广播变量减少大表的 Shuffle;2. 使用缓存(cache/persist)复用中间结果;3. 调整任务调度策略(如 FIFO/FAIR),确保重要作业优先获取资源;4. 监控和分析 Spark UI 中的 DAG 和 Stage 页面,定位瓶颈,如数据倾斜、长尾任务,并针对性优化。
- 【关键点 1】PySpark 通过 DAG 调度器将作业划分为阶段,阶段内任务并行执行。
- 【关键点 2】优化调度核心是减少 Shuffle 和提升数据本地性。
- 【关键点 3】合理设置并行度和资源能显著提升执行效率。
- 【易错点 1】仅增加分区数可能增加调度开销,并非越多越好。
- 【易错点 2】忽视数据倾斜可能导致少数任务执行时间过长。
- 【易错点 3】缓存使用不当可能导致内存溢出(OOM)。