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
相关产品推荐
相关产品推荐

