如何在Spark中并行运行多个作业?基于Java处理YAML文件场景
Java+Spark多YAML任务调度与共享集群资源控制实现方案
实现逻辑说明
- 前4个独立YAML处理任务通过Java
CompletableFuture多线程异步提交到Spark集群,实现并行运行 - 第5个YAML的处理通过
CompletableFuture.allOf()等待前4个任务全部执行结束后再触发,满足依赖要求 - 单应用资源上限通过Spark内置配置参数直接限制,配合集群调度策略避免抢占其他应用资源
依赖配置(Maven示例)
<dependencies> <!-- Spark Core依赖,版本与集群保持一致即可 --> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.12</artifactId> <version>3.3.0</version> <scope>provided</scope> </dependency> <!-- YAML解析工具 --> <dependency> <groupId>org.yaml</groupId> <artifactId>snakeyaml</artifactId> <version>1.33</version> </dependency> </dependencies>
核心代码实现
import org.apache.spark.SparkConf; import org.apache.spark.sql.SparkSession; import org.yaml.snakeyaml.Yaml; import java.io.FileInputStream; import java.util.ArrayList; import java.util.List; import java.util.concurrent.CompletableFuture; public class YamlProcessJob { // 通用YAML处理方法 private static void processSingleYaml(SparkSession spark, String yamlFilePath) { try { // 加载解析YAML配置 Yaml yaml = new Yaml(); FileInputStream fileInputStream = new FileInputStream(yamlFilePath); Object config = yaml.load(fileInputStream); // 此处写入自定义业务处理逻辑,基于config参数调用Spark API完成数据处理 // 处理完成后将结果落盘到指定路径,确保后续任务可读 System.out.printf("%s 处理完成%n", yamlFilePath); } catch (Exception e) { throw new RuntimeException("处理YAML文件失败:" + yamlFilePath, e); } } public static void main(String[] args) { // 初始化Spark配置,直接在配置中限制单应用资源上限 SparkConf conf = new SparkConf() .setAppName("MultiYamlProcessJob") // 限制整个应用最多占用8个CPU核心,可根据集群实际情况调整 .set("spark.cores.max", "8") // 限制单个Executor的内存为4G .set("spark.executor.memory", "4g") // 限制单个Executor占用2个CPU核心 .set("spark.executor.cores", "2") // 若集群启用公平调度器,可指定所属调度池,配合集群侧资源分配规则 .set("spark.scheduler.pool", "normal_pool"); SparkSession spark = SparkSession.builder().config(conf).getOrCreate(); // 前4个独立YAML文件列表 List<String> preTaskYamls = List.of("1.yaml", "2.yaml", "3.yaml", "4.yaml"); List<CompletableFuture<Void>> preTaskFutures = new ArrayList<>(); // 异步提交前4个并行任务 for (String yamlPath : preTaskYamls) { CompletableFuture<Void> taskFuture = CompletableFuture.runAsync( () -> processSingleYaml(spark, yamlPath) ); preTaskFutures.add(taskFuture); } // 阻塞等待所有前置任务执行完成 CompletableFuture.allOf(preTaskFutures.toArray(new CompletableFuture[0])).join(); // 执行依赖前置结果的第5个YAML处理任务 processSingleYaml(spark, "5.yaml"); spark.stop(); } }
资源控制补充说明
spark.cores.max是限制单应用资源的核心参数,不管是Standalone还是YARN集群都生效,调整该参数即可控制当前应用的最大CPU占用上限- 若集群使用YARN作为资源管理器,提交作业时可通过
--queue <队列名>参数指定资源队列,配合YARN队列的资源上限配置,实现双重资源管控 - 集群侧建议启用公平调度器(Fair Scheduler),可保证不同应用按权重分配资源,避免长作业长时间占用所有集群资源
注意事项
- 前4个并行任务的输出路径必须独立配置,避免出现文件写入冲突
- 前置任务的处理结果必须落盘持久化,不要仅保存在缓存中,避免任务异常重试时数据丢失
- SparkSession支持多线程提交作业,无需额外加锁,每个任务的处理逻辑独立构造数据集即可
内容的提问来源于stack exchange,提问作者Rushikesh
相关产品推荐
相关产品推荐

