Spark应用FAIR调度并行执行耗时异常问题求助
让我们一步步拆解你遇到的问题——明明启用了FAIR调度池,GUI也显示任务并行执行,但总耗时和CPU利用率却和FIFO模式完全一致,大概率是几个关键配置或代码逻辑没到位:
1. 线程本地属性未正确传递到并行任务
Spark的spark.scheduler.pool是线程本地属性,而Java parallelStream()使用的ForkJoinPool线程不会自动继承主线程的属性。你在主线程设置了调度池属性,但forEach里的任务是在独立的ForkJoin线程中执行的,这些线程并没有继承到filters池的配置,导致所有查询还是默认提交到default池(FIFO模式),看起来是并行,实际还是串行执行。
修复方案
在每个并行任务内部设置调度池属性,并用try-finally确保属性被重置:
fields.parallelStream().forEach((ColumnMetadata field) -> { SparkContext sc = sqlContext.sparkContext(); try { sc.setLocalProperty("spark.scheduler.pool", "filters"); // 注意:原SQL未过滤字段,每个查询都扫描全表,这里需添加字段过滤减少计算量 Dataset<Row> temp = sqlContext.sql("select distinct tenant_id, user_domain, cube_name, '" + field.getName() + "' as field, value from filter_temp where field = '" + field.getName() + "'"); saveDataFrameToMySQL("analytics_cubes_filters", temp, SaveMode.Append); } finally { sc.setLocalProperty("spark.scheduler.pool", null); } });
2. 数据缓存未真正生效
你调用了dsCube.persist(StorageLevel.MEMORY_ONLY()),但Spark的缓存是懒执行的,只有触发action操作(比如count())才会将数据写入内存。如果没有触发缓存,每个查询都会重新计算filter_temp表,导致重复扫描数据,并行执行也无法提速。
修复方案
在创建临时视图前触发缓存动作:
dsCube.persist(StorageLevel.MEMORY_ONLY()); dsCube.count(); // 强制触发缓存写入 dsCube.createOrReplaceTempView("filter_temp");
3. 调度池配置不合理
你的filters池配置存在两个问题:
weight:1000:如果没有其他调度池,这个值完全没有意义(weight是多池之间的资源分配比例)。minShare:0:minShare是调度器为该池保留的最小资源槽数,设为0意味着调度器不会主动为这个池分配资源,即使有空闲槽位也可能闲置。
优化配置
修改pools.xml,设置与总核数匹配的minShare:
<pool name="filters"> <schedulingMode>FAIR</schedulingMode> <weight>1</weight> <!-- 单池场景下weight任意 --> <minShare>4</minShare> <!-- 等于你的应用总核数,确保池能拿到所有资源 --> </pool>
同时确认spark.scheduler.allocation.file的路径正确,可查看Driver日志确认是否加载到该配置文件。
4. Executor资源与SQL并行度不匹配
你的应用配置是4核10GB内存,但如果SQL查询的并行度设置不合理,会导致资源无法充分利用:
- 默认
spark.sql.shuffle.partitions=200,远大于4核,会导致大量小任务排队,无法并行执行。 - 如果每个查询的任务数太少(比如1个),4核只能同时跑4个任务,但每个查询占1核,总耗时应该接近6秒,而你总耗时24秒,说明可能每个查询都需要独占4核,只能串行执行。
优化配置
- 调整SQL并行度与核数匹配:
spark.sql.shuffle.partitions=4 - 如果是集群模式,确保
spark.executor.instances和spark.executor.cores的配置能提供足够的并行资源(比如2个Executor,每个2核)。
5. MySQL写入成为瓶颈
如果计算部分已经并行,但写入MySQL时因为连接数限制、写入锁或MySQL性能瓶颈,导致串行写入,总耗时还是和串行一致:
- 检查
saveDataFrameToMySQL是否使用了足够大的连接池。 - 确认MySQL的
max_connections配置足够,且没有表级锁导致写入排队。 - 尝试开启Spark的Arrow优化:
spark.sql.execution.arrow.enabled=true提升写入效率。
内容的提问来源于stack exchange,提问作者José Carlos Guevara Turruelles

