Spark作业序列调度及并发控制问题咨询
Spark作业调度相关问题解答
1. 单序列作业的顺序执行(前一个完成再运行下一个)
Spark本身没有直接的跨作业级内置调度机制来自动串联多个独立提交的作业,但可以通过两种内置相关的方式实现,无需依赖外部脚本:
- 同一Spark应用内串联执行:如果这些作业可以放在同一个SparkSession(Spark 2.x+)中运行,直接在代码里按顺序调用action算子即可。Spark的action算子是阻塞式的,前一个action执行完成后才会执行下一行代码,天然保证顺序。示例代码:
val spark = SparkSession.builder().appName("SequentialJobs").getOrCreate() // 作业1:执行action触发计算 spark.read.csv("data1.csv").write.parquet("output1") // 作业2:会在作业1完成后启动 spark.read.parquet("output1").filter("value > 10").write.csv("output2") spark.stop() - 若必须独立提交作业:Spark原生没有内置串联机制,还是需要脚本或外部调度工具,但同应用内的串行逻辑属于Spark内置执行规则,是最直接的实现方式。
2. 多序列作业的并发+同序列串行
Spark原生调度器(FIFO/FAIR)无法直接实现这种分组串行、组间并发的需求,但可以通过以下方式实现:
- 独立Spark应用分组:把A→B→C放在一个Spark应用内按顺序执行,D→E→F放在另一个独立的Spark应用内按顺序执行,同时提交这两个应用。集群调度器(YARN/Standalone)会负责两个应用的并发执行,每个应用内部的作业天然串行,刚好匹配需求。
- FAIR调度器的池机制:如果是同一个Spark应用内的多序列作业,可以给每个序列分配独立的FAIR调度池,池内用FIFO调度,池之间允许并发。步骤如下:
- 修改配置:设置
spark.scheduler.mode=FAIR,并通过fairscheduler.xml定义池属性:<allocations> <pool name="sequence1"> <schedulingMode>FIFO</schedulingMode> <weight>1</weight> </pool> <pool name="sequence2"> <schedulingMode>FIFO</schedulingMode> <weight>1</weight> </pool> </allocations> - 代码中指定作业所属池:
// 序列1的作业A、B使用sequence1池,自动串行 spark.sparkContext.setLocalProperty("spark.scheduler.pool", "sequence1") spark.read.csv("dataA.csv").write.parquet("outputA") spark.read.parquet("outputA").write.csv("outputB") // 序列2的作业D会和sequence1的作业并发执行 spark.sparkContext.setLocalProperty("spark.scheduler.pool", "sequence2") spark.read.csv("dataD.csv").write.parquet("outputD")
- 修改配置:设置
YARN队列机制的适用性
YARN队列可以部分满足需求,但需结合场景配置:
- 对于多序列作业的并发需求:可以将不同序列的作业提交到不同YARN队列,通过队列的资源配置实现组间并发和资源隔离。
- 对于同序列内的串行需求:YARN队列本身不保证队列内作业的执行顺序,必须配合Spark应用内的串行逻辑(如同一应用内按顺序执行),或外部脚本控制提交顺序,才能实现同序列作业的串行。
内容的提问来源于stack exchange,提问作者maxwellray
相关产品推荐
相关产品推荐

