如何根据输入文件大小调整Spark executor核心、内存及实例数量
Spark针对输入文件大小动态调整Executor配置方案
基础参数匹配逻辑
你可以根据输入数据规模先计算出对应参数值,再传入作业,核心计算逻辑如下:
- Executor核心数:通用场景下固定配置为45核即可,这个数值的CPU上下文切换开销最低,利用率最优,特殊重计算场景可以下调到23核
- 单Executor内存:和核心数匹配即可,1核对应24GB内存,同时预留10%15%的内存作为系统开销,比如4核Executor配8~16GB内存即可,
spark.executor.memoryOverhead可以单独配置,也可以让Spark按默认10%的比例自动计算 - Executor实例数:根据输入数据的总分片数计算,Spark默认HDFS分片大小为128MB,总分片数=总输入大小/128MB,不足1片按1片计算;单个Executor的并行task数等于其核心数,所以Executor实例数=总分片数/Executor核心数,再预留20%左右的余量即可,不要超过集群的最大可用资源上限
具体实现方式
脚本提交场景
你可以在提交作业的Shell/Python脚本里先统计输入文件的大小,动态拼接spark-submit的参数即可,Shell示例如下:
# 输入路径 INPUT_PATH="hdfs://xxx/your/input/path" # 统计输入总大小,单位GB TOTAL_SIZE=$(hdfs dfs -du -s $INPUT_PATH | awk '{print $1/1024/1024/1024}') # 固定配置 EXECUTOR_CORES=4 MEM_PER_CORE=3 MAX_EXECUTORS=100 # 集群最多可给该作业分配的Executor数量 # 计算总分片数,默认分片128MB TOTAL_SPLITS=$(echo "scale=0; $TOTAL_SIZE / 0.128 + 1" | bc) # 计算Executor实例数,留20%余量 EXECUTOR_INSTANCES=$(echo "scale=0; $TOTAL_SPLITS / $EXECUTOR_CORES * 1.2 + 1" | bc) # 不超过集群上限 EXECUTOR_INSTANCES=$(( EXECUTOR_INSTANCES > MAX_EXECUTORS ? MAX_EXECUTORS : EXECUTOR_INSTANCES )) # 计算Executor内存 EXECUTOR_MEM="${EXECUTOR_CORES * MEM_PER_CORE}g" # 动态提交作业 spark-submit \ --conf spark.executor.cores=$EXECUTOR_CORES \ --conf spark.executor.memory=$EXECUTOR_MEM \ --conf spark.executor.instances=$EXECUTOR_INSTANCES \ --class your.main.class.Name \ your-app.jar $INPUT_PATH
代码内部动态配置场景
你也可以在Spark代码初始化SparkContext/SparkSession之前,先统计输入文件大小,动态设置配置,Scala示例如下:
import org.apache.hadoop.fs.{FileSystem, Path} import org.apache.spark.sql.SparkSession object YourSparkApp { def main(args: Array[String]): Unit = { val inputPath = args(0) // 统计输入文件总大小,单位GB val hadoopConf = new org.apache.hadoop.conf.Configuration() val fs = FileSystem.get(hadoopConf) val totalSizeByte = fs.getContentSummary(new Path(inputPath)).getLength val totalSizeGB = totalSizeByte / (1024 * 1024 * 1024.0) // 固定配置 val executorCores = 4 val memPerCore = 3 val maxExecutors = 100 // 计算参数 val totalSplits = Math.ceil(totalSizeGB / 0.128).toInt var executorInstances = Math.ceil(totalSplits / executorCores * 1.2).toInt executorInstances = Math.min(executorInstances, maxExecutors) val executorMem = s"${executorCores * memPerCore}g" // 初始化SparkSession val spark = SparkSession.builder() .config("spark.executor.cores", executorCores) .config("spark.executor.memory", executorMem) .config("spark.executor.instances", executorInstances) // 其他配置 .appName("DynamicConfigDemo") .getOrCreate() // 后续业务逻辑 } }
可选方案:开启集群动态资源分配
如果你不想自己写逻辑计算参数,也可以直接开启Spark的动态资源分配功能,让Spark根据作业的实际task数量自动扩缩容Executor,不需要提前固定参数,配置示例如下:
spark.dynamicAllocation.enabled=true spark.dynamicAllocation.minExecutors=1 spark.dynamicAllocation.maxExecutors=100 # 调整为集群允许的最大值 spark.dynamicAllocation.executorIdleTimeout=60s # 空闲60秒就释放Executor
该方案适合作业逻辑多变、输入大小波动大的场景,不用手动维护参数计算逻辑。
边界调整建议
- 如果输入路径下小文件极多,实际文件数量远大于计算出来的总分片数,就用实际文件数量代替总分片数计算Executor实例数,避免单Task处理过多小文件拖慢性能
- 如果作业包含大量join、聚合、缓存等重内存消耗的逻辑,可以把内存配比提升到1核对应46GB内存,Executor实例数也额外增加20%30%的余量
内容的提问来源于stack exchange,提问作者Nikhil Thomas
相关产品推荐
相关产品推荐

