You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Spark中并行运行多个作业?基于Java处理YAML文件场景

Java+Spark多YAML任务调度与共享集群资源控制实现方案

实现逻辑说明

  • 前4个独立YAML处理任务通过JavaCompletableFuture多线程异步提交到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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 02:18:04