前些天发现了一个巨牛的人工智能学习网站,通俗易懂,风趣幽默,忍不住分享一下给大家。点击跳转到网站:https://www.captainai.net/dongkelun
前言
在第二篇文章 Spark 源码 | SparkSubmit 提交流程分析(二) 中,我们分析了 SparkSubmit 的提交流程,其中 doSubmit 方法会调用 parseArguments 解析命令行参数。本文详细介绍参数解析的过程。
版本
Spark 3.2.3
SparkSubmitArguments 概述
SparkSubmitArguments 继承自 SparkSubmitArgumentsParser,用于解析 spark-submit 命令行参数。
1// SparkSubmitArguments.scala 2private[deploy] class SparkSubmitArguments(args: Seq[String], env: Map[String, String] = sys.env) 3 extends SparkSubmitArgumentsParser with Logging { 4
字段定义
1// --master MASTER_URL 2// Spark Master 地址,如 spark://host:port, mesos://host:port, yarn, k8s://https://host:port, local(默认 local[*]) 3var master: String = null 4 5// --deploy-mode DEPLOY_MODE 6// 部署模式,client 表示在本地启动 Driver,cluster 表示在集群的 Worker 节点上启动 Driver(默认 client) 7var deployMode: String = null 8 9// --executor-memory MEM 10// 每个 Executor 的内存(如 1000M, 2G)(默认 1G) 11var executorMemory: String = null 12 13// --executor-cores NUM 14// 每个 Executor 的 CPU 核心数(YARN 和 K8S 模式默认 1,Standalone 模式默认为 Worker 上的所有核心) 15var executorCores: String = null 16 17// --total-executor-cores NUM 18// 所有 Executor 的总 CPU 核心数 19var totalExecutorCores: String = null 20 21// --properties-file FILE 22// 配置文件路径,用于加载额外的 Spark 配置。如果未指定,则默认查找 conf/spark-defaults.conf 23var propertiesFile: String = null 24 25// --driver-memory MEM 26// Driver 的内存(如 1000M, 2G) 27var driverMemory: String = null 28 29// --driver-class-path 30// 额外添加到 Driver classpath 的路径。注意:通过 --jars 添加的 jar 包会自动包含在 classpath 中 31var driverExtraClassPath: String = null 32 33// --driver-library-path 34// 额外添加到 Driver 库路径的环境变量 35var driverExtraLibraryPath: String = null 36 37// --driver-java-options 38// 传递给 Driver 的额外 Java 选项 39var driverExtraJavaOptions: String = null 40 41// --queue QUEUE_NAME 42// YARN 队列名称(默认 "default") 43var queue: String = null 44 45// --num-executors NUM 46// Executor 的数量(默认 2)。启用动态分配时,实际初始数量取 --num-executors、 47// spark.dynamicAllocation.minExecutors 和 spark.dynamicAllocation.initialExecutors 三者中的最大值 48var numExecutors: String = null 49 50// --files FILES 51// 逗号分隔的文件列表,会被放置在每个 Executor 的工作目录中。在 Executor 中可通过 SparkFiles.get(fileName) 获取文件路径 52var files: String = null 53 54// --archives ARCHIVES 55// 逗号分隔的归档文件列表,会被提取到每个 Executor 的工作目录中 56var archives: String = null 57 58// --class CLASS_NAME 59// 应用程序的主类(用于 Java / Scala 应用) 60var mainClass: String = null 61 62// <app jar | python file | R file> 63// 第一个未被识别的选项被视为"主资源"(primary resource) 64var primaryResource: String = null 65 66// --name NAME 67// 应用程序的名称 68var name: String = null 69 70// 应用程序参数 71var childArgs: ArrayBuffer[String] = new ArrayBuffer[String]() 72 73// --jars JARS 74// 逗号分隔的 JAR 包列表,会被添加到 Driver 和 Executor 的 classpath 中 75var jars: String = null 76 77// --packages 78// 逗号分隔的 Maven 坐标列表,用于添加 JAR 包到 Driver 和 Executor 的 classpath。 79// 会先搜索本地 Maven 仓库,然后是 Maven 中央仓库,以及 --repositories 指定的远程仓库。 80// 坐标格式为 groupId:artifactId:version 81var packages: String = null 82 83// --repositories 84// 逗号分隔的远程仓库地址,用于搜索 --packages 指定的 Maven 坐标 85var repositories: String = null 86 87var ivyRepoPath: String = null 88 89var ivySettingsPath: Option[String] = None 90 91// --exclude-packages 92// 逗号分隔的 groupId:artifactId 列表,在解析 --packages 提供的依赖时排除这些包,以避免依赖冲突 93var packagesExclusions: String = null 94 95// --verbose, -v 96// 打印额外的调试输出 97var verbose: Boolean = false 98 99var isPython: Boolean = false 100 101// --py-files PY_FILES 102// 逗号分隔的 .zip、.egg 或 .py 文件列表,会被添加到 Python 应用的 PYTHONPATH 中 103var pyFiles: String = null 104 105var isR: Boolean = false 106 107var action: SparkSubmitAction = null 108 109// 通过 `--conf` 和配置文件加载的 Spark 配置属性 110val sparkProperties: HashMap[String, String] = new HashMap[String, String]() 111 112// --proxy-user NAME 113// 提交应用程序时模拟的用户名。不能与 --principal / --keytab 同时使用 114var proxyUser: String = null 115 116// --principal PRINCIPAL 117// 用于登录 KDC 的 Principal(仅 YARN/K8s 模式) 118var principal: String = null 119 120// --keytab KEYTAB 121// 包含上述 Principal 对应 keytab 的文件路径 122var keytab: String = null 123 124private var dynamicAllocationEnabled: Boolean = false 125 126// --supervise 127// 如果设置,Driver 失败时会自动重启(仅 Spark Standalone 或 Mesos 的 cluster 部署模式) 128var supervise: Boolean = false 129 130// --driver-cores NUM 131// Driver 使用的 CPU 核心数,仅在 cluster 模式可用(默认 1) 132var driverCores: String = null 133 134// --kill SUBMISSION_ID 135// 如果设置,杀死指定的 Driver(Spark Standalone、Mesos 或 K8s 的 cluster 部署模式) 136var submissionToKill: String = null 137 138// --status SUBMISSION_ID 139// 如果设置,查询指定 Driver 的状态 140var submissionToRequestStatusFor: String = null 141 142// 内部使用,用于 REST API 143var useRest: Boolean = false 144
解析流程
1. 构造函数调用顺序
1// Set parameters from command line arguments 2// 从命令行参数设置参数 3parse(args.asJava) 4 5// Populate `sparkProperties` map from properties file 6// 从属性文件填充 sparkProperties 7mergeDefaultSparkProperties() 8// Remove keys that don't start with "spark." from `sparkProperties`. 9// 从 sparkProperties 中移除不以 "spark." 开头的键 10ignoreNonSparkProperties() 11// Use `sparkProperties` map along with env vars to fill in any missing parameters 12// 使用 sparkProperties 和环境变量填充缺失的参数 13loadEnvironmentArguments() 14 15useRest = sparkProperties.getOrElse("spark.master.rest.enabled", "false").toBoolean 16 17validateArguments() 18
解析顺序:
- parse(args.asJava) - 从命令行参数设置参数
- mergeDefaultSparkProperties() - 从属性文件填充 sparkProperties
- ignoreNonSparkProperties() - 从 sparkProperties 中移除不以 “spark.” 开头的键
- loadEnvironmentArguments() - 使用 sparkProperties 和环境变量填充缺失的参数
- validateArguments() - 验证参数
2. 解析命令行参数
继承自 SparkSubmitArgumentsParser(继承自 SparkSubmitOptionParser),重写 handle、handleUnknown、handleExtraArgs 三个回调方法处理参数。
SparkSubmitOptionParser parse 方法解析逻辑
parse(List<String> args) 方法的遍历逻辑:
1// 遍历所有参数 2for (idx = 0; idx < args.size(); idx++) { 3 // 支持 --master=spark://xxx 格式(等号分隔) 4 // 检查是否是 --xxx=value 格式,若是则拆分 5 Matcher m = eqSeparatedOpt.matcher(arg); 6 if (m.matches()) { 7 arg = m.group(1); 8 value = m.group(2); 9 } 10 11 // Look for options with a value. 12 // 第一步:在 opts 数组中查找(带参数的选项,如 --master, --class 等) 13 String name = findCliOption(arg, opts); 14 if (name != null) { 15 // If the option requires a value, fetch it. 16 // 命令行格式:--选项 值(分开的两参数),如 --master yarn 17 if (value == null) { 18 if (idx == args.size() - 1) { 19 throw new IllegalArgumentException( 20 String.format("Missing argument for option '%s'.", arg)); 21 } 22 idx++; 23 value = args.get(idx); 24 } 25 if (!handle(name, value)) { 26 // handle 返回 false 的情况:--version(--help 会直接退出,不会返回) 27 break; 28 } 29 continue; 30 } 31 32 // Look for a switch. 33 // 第二步:在 switches 数组中查找(不带参数的选项,如 --help, --verbose 等) 34 name = findCliOption(arg, switches); 35 if (name != null) { 36 // handle 返回 false 的情况:--version(打印版本后退出) 37 if (!handle(name, null)) { 38 break; 39 } 40 continue; 41 } 42 43 // 第三步:既不在 opts 也不在 switches 中,调用 handleUnknown 处理未知选项 44 // SparkSubmitArguments 中将第一个未知选项作为 primaryResource,handleUnknown 始终返回 false 45 // 返回 false 则 break 退出循环,剩余参数作为应用参数 46 if (!handleUnknown(arg)) { 47 break; 48 } 49} 50 51// 跳过第一个未知选项(primaryResource),只把剩余参数传给 handleExtraArgs 52// 如果正常遍历完(无 break),idx == args.size(),不需要 idx++ 53// 如果因 break 提前退出,idx < args.size(),需要跳过 primaryResource 54if (idx < args.size()) { 55 idx++; 56} 57 58// 第四步:解析完成后,剩余的参数作为应用参数调用 handleExtraArgs 59// 注意:循环可能因 break 提前退出,导致 idx < args.size() 60// 如:handleUnknown 返回 false 时(SparkSubmitArguments 将第一个未知选项作为 primaryResource) 61// spark-submit 命令格式:spark-submit [options] <app jar | python file> [app arguments] 62// 剩余的参数即 [app arguments],如 arg1 arg2 arg3 63handleExtraArgs(args.subList(idx, args.size())); 64
SparkSubmitArguments 重写的方法
1. handle 方法
处理已知的带参数选项,根据 opt(选项名)将 value 赋值给对应的字段。其中 CONF case 会直接将 --conf 的值存入 sparkProperties:
1override protected def handle(opt: String, value: String): Boolean = { 2 opt match { 3 case NAME => 4 name = value 5 case MASTER => 6 master = value 7 case CLASS => 8 mainClass = value 9 case DEPLOY_MODE => 10 if (value != "client" && value != "cluster") { 11 error("--deploy-mode must be either \"client\" or \"cluster\"") 12 } 13 deployMode = value 14 case NUM_EXECUTORS => 15 numExecutors = value 16 // ... 其他参数 17 // `--conf` KEY=VALUE 格式的配置项,直接存入 sparkProperties 18 // 由于是直接赋值(覆盖写),后续合并配置文件时,`--conf` 的值会被保留 19 case CONF => 20 val (confName, confValue) = SparkSubmitUtils.parseSparkConfProperty(value) 21 sparkProperties(confName) = confValue 22 // ... 其他参数 23 } 24 // --version 时返回 false,其他情况返回 true 25 action != SparkSubmitAction.PRINT_VERSION 26} 27
2. handleUnknown 方法
处理未知选项,将第一个未知选项作为 primaryResource(主资源),如 myapp.jar、spark-shell、hdfs:///path/to/app.jar:
1override protected def handleUnknown(opt: String): Boolean = { 2 // 如果选项以 "-" 开头但不在已知选项中,报错 3 if (opt.startsWith("-")) { 4 error(s"Unrecognized option '$opt'.") 5 } 6 7 // 设置 primaryResource(主资源),如 myapp.jar、spark-shell、hdfs:///path/to/app.jar 8 // 如果不是 shell(spark-shell、pyspark-shell、sparkr-shell)或内部选项(spark-internal),则通过 Utils.resolveURI 解析为标准 URI 格式 9 // 如:./app.jar -> file:/path/to/app.jar,/path/app.jar -> file:/path/app.jar,hdfs://xxx -> hdfs://xxx 10 // 否则直接作为字符串(如 spark-shell 用于交互式界面,spark-internal 表示没有指定应用程序资源,用于交互式工具如 spark-sql 或测试场景) 11 primaryResource = 12 if (!SparkSubmit.isShell(opt) && !SparkSubmit.isInternal(opt)) { 13 Utils.resolveURI(opt).toString 14 } else { 15 opt 16 } 17 18 // 判断是否为 Python 或 R 应用 19 isPython = SparkSubmit.isPython(opt) 20 isR = SparkSubmit.isR(opt) 21 22 // 始终返回 false,退出循环 23 false 24} 25
3. handleExtraArgs 方法
将解析剩余的应用参数添加到 childArgs:
1override protected def handleExtraArgs(extra: JList[String]): Unit = { 2 // 在 SparkSubmitArguments 类中只有这一处赋值,没看到有其它地方赋值,感觉两种写法效果可能是相同的 3 // 使用 ++= 追加而非 = 赋值,可能为了和 SparkSubmit.scala 中对 childArgs 的操作风格保持一致 4 // childArgs 是 ArrayBuffer[String],extra.asScala 是 Seq[String] 5 // 使用 ++= 避免了显式类型转换(如 childArgs = extra.asScala.to[ArrayBuffer]) 6 childArgs ++= extra.asScala 7} 8
3. 合并默认属性文件
将配置文件(默认 conf/spark-defaults.conf)中的属性合并到 sparkProperties 中。如果用户通过 --conf 指定了相同的配置项,则保留 --conf 的值(优先级更高)。
1private def mergeDefaultSparkProperties(): Unit = { 2 // Use common defaults file, if not specified by user 3 // 如果用户未指定,则使用默认属性文件 4 propertiesFile = Option(propertiesFile).getOrElse(Utils.getDefaultPropertiesFile(env)) 5 // Honor `--conf` before the defaults file 6 // `--conf` 的优先级高于默认属性文件 7 defaultSparkProperties.foreach { case (k, v) => 8 if (!sparkProperties.contains(k)) { 9 sparkProperties(k) = v 10 } 11 } 12} 13
4. 过滤非 Spark 属性
将 sparkProperties 中不以 spark. 开头的键移除掉。sparkProperties 中的值来自 --conf 和配置文件,不包含环境变量等来源。
1private def ignoreNonSparkProperties(): Unit = { 2 sparkProperties.keys.foreach { k => 3 if (!k.startsWith("spark.")) { 4 sparkProperties -= k 5 logWarning(s"Ignoring non-Spark config property: $k") 6 } 7 } 8} 9
5. 加载环境变量参数
将命令行未指定的参数从环境变量和 sparkProperties 中补充进来。例如:master 未在命令行指定时,会依次从 spark.master 配置、环境变量 MASTER 中获取。
1private def loadEnvironmentArguments(): Unit = { 2 // Load arguments from environment variables, Spark properties etc. 3 // 从环境变量、Spark 属性等加载参数 4 master = Option(master) 5 .orElse(sparkProperties.get("spark.master")) 6 .orElse(env.get("MASTER")) 7 .orNull 8 // ... 其他参数类似 9} 10
优先级:命令行参数 > sparkProperties > 环境变量 > 默认值
6. 参数验证
验证必填参数是否存在以及参数值是否合法。例如:必须指定 primaryResource(主资源),内存/核心数必须为正数,YARN 模式必须设置 HADOOP_CONF_DIR 等。
1private def validateSubmitArguments(): Unit = { 2 // Ensure that required fields exists. Call this only once all defaults are loaded. 3 // 确保必填字段存在。调用此方法时所有默认值已加载完成 4 if (args.length == 0) { 5 printUsageAndExit(-1) 6 } 7 if (primaryResource == null) { 8 error("Must specify a primary resource (JAR or Python or R file)") 9 } 10 // ... 其他验证 11} 12
参数优先级
| 来源 | 优先级 |
|---|---|
| 命令行参数 | 1(最高) |
| --conf | 2 |
| 配置文件 | 3 |
| 环境变量 | 4 |
| 默认值 | 5(最低) |
总结
SparkSubmitArguments 的解析流程:
- parse - 解析命令行参数
- mergeDefaultSparkProperties - 合并配置文件中的属性
- ignoreNonSparkProperties - 过滤非 spark. 开头的属性
- loadEnvironmentArguments - 用环境变量和 sparkProperties 填充空缺参数
- validateArguments - 验证参数合法性
优先级:命令行 > --conf > 配置文件 > 环境变量 > 默认值
下一次我们将分析 Yarn Client、Yarn Cluster、Standalone Client、Standalone Cluster 模式的详细提交流程。
