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

Java Spark3.0如何编程获取Spark集群可用节点及Executor数量

Spark 3.0 Java版本可编程获取集群资源实现自动分区方案

先明确几个常见踩坑点

  • java.lang.Runtime.getRuntime.availableProcessors 只能获取当前Driver进程所在节点的CPU核数,本地调试模式可用,提交到集群后拿到的值和集群整体可用资源完全无关
  • 旧版本通过反射取BlockManager、调用sc.getExecutorStorageStatus.length - 1计算Executor数的方案在Spark 3.x已经失效:Spark 3调整了BlockManagerMaster的状态返回权限和逻辑,未完成注册时接口只会返回Driver自身的状态,硬写反射破封装的逻辑对小版本兼容性极差,3.0.0和3.0.2的私有字段偏移就可能存在差异,生产环境用很容易随版本升级报错。
  • sc.defaultParallelism 不能直接当总可用核数用:这个值是Spark启动阶段读取静态配置生成的,如果开启动态资源分配,这个值不会随Executor扩缩容更新;如果用户手动配置过该参数,返回值和实际集群资源完全不匹配。

可用方案1:SparkListener监听注册事件(最稳定,全部署模式兼容)

这个方案是Spark原生提供的能力,不需要反射、不依赖私有API,适配Standalone/Yarn/K8s所有部署模式,逻辑是等所有Executor完成向Driver的注册流程后,累加每个Executor的实际分配核数,拿到准确的总可用核心数。
Java实现代码如下:

import org.apache.spark.SparkConf;
import org.apache.spark.SparkContext;
import org.apache.spark.scheduler.SparkListener;
import org.apache.spark.scheduler.SparkListenerExecutorAdded;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

public class AutoPartitionApp {
    public static void main(String[] args) throws InterruptedException {
        SparkConf conf = new SparkConf().setAppName("AutoPartitionDemo");
        SparkContext sc = new SparkContext(conf);

        AtomicInteger totalExecutorCores = new AtomicInteger(0);
        AtomicInteger executorCount = new AtomicInteger(0);
        // 读取静态配置的预期Executor数,适配动态/静态资源分配场景
        int expectedExec = conf.getInt("spark.executor.instances",
                conf.getInt("spark.dynamicAllocation.initialExecutors", 1));
        CountDownLatch registerLatch = new CountDownLatch(expectedExec);

        // 注册监听器,每个Executor上线时累加核数
        sc.addSparkListener(new SparkListener() {
            @Override
            public void onExecutorAdded(SparkListenerExecutorAdded event) {
                totalExecutorCores.addAndGet(event.executorInfo().totalCores());
                executorCount.incrementAndGet();
                registerLatch.countDown();
            }
        });

        // 最多等待10秒完成Executor注册,避免作业无限阻塞
        registerLatch.await(10, TimeUnit.SECONDS);
        // 常规场景分区数设为总核数的2~3倍即可,可按业务IO/计算特性调整
        int optimalPartitionNum = totalExecutorCores.get() * 2;

        // 后续读取数据、重分区时直接传入计算好的分区数即可,例:
        // JavaRDD<String> rdd = sc.textFile("hdfs:///data/path", optimalPartitionNum);

        sc.stop();
    }
}

如果开启动态资源分配,可以去掉CountDownLatch的固定等待逻辑,改为应用启动后等待5~10秒收集已注册的Executor资源,后续再通过监听SparkListenerExecutorAdded、SparkListenerExecutorRemoved事件动态更新总核数,触发作业重分区即可。


可用方案2:StatusTracker直接拉取Executor状态(实现简单,适合测试/固定资源集群)

如果不想写监听器,可以在SparkContext初始化完成后,预留3~5秒等待Executor完成注册,直接从Spark内置的状态追踪器拉取活跃Executor信息:

SparkContext sc = new SparkContext(conf);
// 等待Executor注册,可按集群启动速度调整等待时长
TimeUnit.SECONDS.sleep(5);

int totalCores = 0;
int executorNum = 0;
for (String execId : sc.statusTracker().getExecutorInfos()) {
    // 排除Driver自身的资源
    if ("driver".equals(execId)) continue;
    totalCores += sc.statusTracker().getExecutorInfo(execId).get().totalCores();
    executorNum++;
}
int optimalPartitionNum = totalCores * 2;

这个方案的缺点是固定等待时长无法适配集群资源不足、Executor启动慢的场景,可能出现资源统计不全的问题,适合资源固定的生产集群或者本地测试场景用。


内容的提问来源于stack exchange,提问作者capt-mac

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 19:27:20