Apache DolphinScheduler Spark 任务节点完全指南:spark-submit 与 Spark SQL 提交原理与实战
【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/GitHub_Trending/dol/dolphinscheduler
Spark 任务节点是 Apache DolphinScheduler 工作流中用于提交 Apache Spark 应用程序的核心任务类型,它屏蔽了手工拼接spark-submit/spark-sql命令的繁琐过程,通过表单化参数即可把任意 Java、Scala、Python 或 SQL 的 Spark 作业纳入低代码编排的 DAG 中。本文以 Spark 任务节点官方文档 为骨架,结合仓库中dolphinscheduler-task-spark插件的源码与测试用例,完整讲解任务创建步骤、全部参数含义、参数到底层命令的映射规则,并给出 WordCount 与 Spark SQL 两个可直接复用的实战案例。读完本文,你将能够独立配置并排产一个生产可用的 Spark 任务节点,同时理解其底层命令组装逻辑,便于排查提交类问题。
Spark 任务节点概述
Spark 任务节点用于执行 Spark 应用程序。当任务被调度到 Worker 执行时,Worker 会依据用户配置,通过以下两种方式之一向 Spark 集群提交作业:
spark-submit方式:适用于 Java / Scala / Python 类型的 Spark 应用,向集群提交主程序(JAR 包或 Python 文件)并附带启动参数。spark-sql方式:适用于 SQL 类型的作业,以spark-sql -f <filename>的形式执行 SQL 脚本。
在源码层面,这两种方式的命令入口定义在 SparkConstants.java 中:
public static final String SPARK_SQL_COMMAND = "${SPARK_HOME}/bin/spark-sql"; public static final String SPARK_SUBMIT_COMMAND = "${SPARK_HOME}/bin/spark-submit";命令中的${SPARK_HOME}占位符由执行环境提供,意味着 Worker 所在机器必须安装 Spark 并正确配置SPARK_HOME(详见下文"环境准备"一节)。
SparkTask.java 中的getScript()方法负责选择命令:当ProgramType为SQL时走spark-sql,否则走spark-submit,随后把populateSparkOptions()组装出的参数拼接为最终命令字符串。该逻辑与 ProgramType.java 中定义的JAVA / SCALA / PYTHON / SQL四种程序类型一一对应。
创建 Spark 任务
在 DolphinScheduler Web UI 中按以下步骤创建 Spark 任务节点:
- 进入
项目管理 -> 项目名称 -> 工作流定义,点击创建工作流按钮进入 DAG 编辑页面; - 从左侧工具栏拖拽SPARK节点到画布中;
- 双击节点,在右侧表单中按下文"任务参数详解"填写配置;
- 保存节点并完成 DAG 连线,即可作为工作流的一部分被调度执行。
任务参数详解
Spark 任务节点的参数由两部分构成:一部分是所有任务类型共享的默认参数,另一部分是 Spark 任务专属参数。
默认任务参数
节点名称、运行标志、描述、任务优先级、Worker 分组、任务组名称、环境名称、失败重试次数、失败重试间隔、CPU 配额、最大内存、超时告警、延时执行时间、资源、前置任务等通用参数,请参考 DolphinScheduler Task Parameters Appendix 中Default Task Parameters一节。其中环境名称参数在配置了 Spark 相关环境(见下文)后,可用于为任务指定独立的执行环境。
Spark 专属参数
| 参数 | 说明 |
|---|---|
| 程序类型 (Program type) | 支持 Java、Scala、Python 和 SQL 四种类型。 |
| 主函数的类 (The class of main function) | Spark 程序入口 Main Class 的全限定名(full path),例如org.apache.spark.examples.JavaWordCount。仅 Java / Scala 类型需要。 |
| Master | 集群的 Master URL,例如yarn、spark://localhost:7077或 Kubernetes 集群地址。留空时由插件按规则推导(见"命令组装"一节)。 |
| 主程序包 (Main jar package) | 通过资源中心上传的 Spark 主程序 JAR 包(Java / Scala 使用)或 Python 文件(Python 类型使用)。 |
| SQL 脚本 (SQL scripts) | Spark SQL 执行的.sql文件(SQL 类型使用),从资源中心选择。 |
| 部署方式 (Deployment mode) | spark-submit支持cluster、client、local三种模式;spark-sql支持client和local两种模式。 |
| 命名空间(集群)(Namespace) | 选择命名空间后,作业提交到原生 Kubernetes 集群;不选择时默认提交到 YARN 集群。 |
| 任务名称 (Task name) | Spark 应用名称,对应spark-submit --name,也是 DAG 中的节点名称。 |
| Driver 核心数 (Driver core number) | 设置 Driver 使用的核心数,对应--conf spark.driver.cores=N,可按实际生产环境调整。 |
| Driver 内存大小 (Driver memory size) | 设置 Driver 内存大小,对应--conf spark.driver.memory=<SIZE>。 |
| Executor 数量 (Number of Executor) | 设置 Executor 个数,对应--conf spark.executor.instances=N。 |
| Executor 内存大小 (Executor memory size) | 设置每个 Executor 的内存大小,对应--conf spark.executor.memory=<SIZE>。 |
| Yarn 队列 (Yarn queue) | 设置提交作业使用的 YARN 队列,默认使用default队列,对应--queue <queue>。 |
| 主程序参数 (Main program parameters) | 传给 Spark 主程序的应用参数(app arguments),支持 DolphinScheduler 自定义参数变量替换。 |
| 可选参数 (Optional parameters) | 追加的 Spark 命令选项,如--jars、--files、--archives、--conf等,会原样拼接到命令末尾。 |
| 资源 (Resource) | 当参数中引用了资源中心文件时,在此处指定关联的资源文件,执行前会自动下载到 Worker 本地。 |
| 自定义参数 (Custom parameter) | Spark 任务局部自定义参数,会将脚本中${variable}形式的内容替换为对应值。 |
| 前置任务 (Predecessor task) | 为当前任务选择前置任务,被选中的任务将成为当前任务的上游节点。 |
上述参数在插件侧由 SparkParameters.java 承载(mainJar、mainClass、master、deployMode、mainArgs、driverCores、driverMemory、numExecutors、executorCores、executorMemory、appName、yarnQueue、others、programType、rawScript、namespace、resourceList、sqlExecutionType等字段),任务保存时通过checkParameters()校验参数合法性:
ProgramType为 SQL 时,rawScript不能为空;ProgramType为 Java / Scala / Python 时,mainJar不能为空;ProgramType本身不能为空。
实战示例一:spark-submit 提交 WordCount 程序
WordCount 是大数据生态中最常见的入门案例,适用于 MapReduce、Flink、Spark 等计算框架,核心目的是统计输入文本中相同单词的出现次数。下面演示如何在 DolphinScheduler 中完整配置一个 Spark WordCount 任务。
第一步:在 DolphinScheduler 中配置 Spark 环境
在生产环境使用 Spark 任务类型前,必须先在执行任务的 Worker 机器上准备好 Spark 运行环境。DolphinScheduler 的任务执行环境配置文件为 script/env/dolphinscheduler_env.sh(部署到安装目录后通常位于bin/env/或script/env/下),需要在其中声明JAVA_HOME、SPARK_HOME等变量。仓库中 CI 使用的完整示例(mysql_with_zookeeper_registry/dolphinscheduler_env.sh)给出了关键配置:
export JAVA_HOME=${JAVA_HOME:-/opt/java/openjdk} export HADOOP_HOME=${HADOOP_HOME:-/opt/soft/hadoop} export HADOOP_CONF_DIR=${HADOOP_CONF_DIR:-/opt/soft/hadoop/etc/hadoop} export SPARK_HOME=${SPARK_HOME:-/opt/soft/spark} export PATH=$HADOOP_HOME/bin:$SPARK_HOME/bin:$JAVA_HOME/bin:$PATH要点说明:
SPARK_HOME必须指向 Worker 机器上已安装的 Spark 发行版目录,${SPARK_HOME}/bin/spark-submit与${SPARK_HOME}/bin/spark-sql依赖该变量定位可执行脚本;- 若提交到 YARN,还需正确配置
HADOOP_HOME与HADOOP_CONF_DIR,确保spark-submit能拿到yarn-site.xml等集群配置; - 该文件在每次任务执行时都会被 source,因此不要在文件中存放数据库密码等敏感信息。
第二步:上传主程序包到资源中心
使用 Spark 任务节点前,需要先将主程序 JAR 包上传到资源中心(Resource Center)。资源中心支持本地文件系统、HDFS、S3、OSS、OBS、COS 等多种存储后端,具体配置方法参见 资源中心配置文档。配置完成后,直接在资源中心页面通过拖拽方式上传目标文件即可。上传成功后,在主程序包参数中即可选中该 JAR。
第三步:配置 Spark 任务节点
根据上文参数表在节点表单中填写以下内容:
- 程序类型:Java(或 Scala,取决于主程序的实现语言);
- 主函数的类:填写 Main Class 的全限定名,如
org.apache.spark.examples.JavaWordCount; - 部署方式:
client或cluster(提交到 YARN 集群时常用cluster); - 主程序包:选择资源中心已上传的 JAR;
- 任务名称:如
spark-wordcount; - Driver / Executor 相关资源配置:按生产环境实际规格填写核心数与内存;
- 主程序参数:传入输入路径、输出路径等应用参数;
- 如作业需要额外依赖,可在可选参数中追加
--jars、--conf等选项。
保存并发布工作流后,DolphinScheduler 会在 Worker 上生成类似如下的提交命令(参考 SparkTaskTest.java 中断言的命令格式):
${SPARK_HOME}/bin/spark-submit --master yarn --deploy-mode client \ --class org.apache.spark.examples.JavaWordCount \ --conf spark.driver.cores=1 --conf spark.driver.memory=512M \ --conf spark.executor.instances=2 --conf spark.executor.cores=2 \ --conf spark.executor.memory=1G --name spark \ /lib/dolphinscheduler-task-spark.jar实战示例二:Spark SQL 执行 DDL 与 DML 语句
Spark SQL 类型的任务适用于直接以 SQL 操作 Spark 数据源的场景。官方文档给出的案例是:创建视图表terms并写入三行数据,再创建一个 parquet 格式的表wc并判断表是否存在,最后将视图表terms中的数据插入wc。
配置要点:
- 程序类型:选择
SQL; - 部署方式:Spark SQL 支持
client与local两种模式,不支持 cluster 模式; - SQL 脚本:从资源中心选择一个
.sql文件作为脚本内容来源。
关于 SQL 内容来源,源码支持两种方式(由sqlExecutionType字段区分,见 SparkConstants.java):
- FILE 模式:从资源中心读取
.sql文件内容,多个文件时默认取第一个(源码会打印告警more than 1 files detected, use the first one by default); - SCRIPT 模式:使用表单中的内联脚本(
rawScript)。
无论哪种方式,插件都会先把 SQL 内容写到 Worker 执行目录下的{taskAppId}_node.sql临时文件,再以spark-sql -f <filename>执行(见 SparkTask.java)。写入前会执行参数替换(replaceParam),因此 SQL 中支持使用${变量}形式的自定义参数。生成的最终命令形如(参考 SparkTaskTest.java):
${SPARK_HOME}/bin/spark-sql --master yarn --deploy-mode client \ --conf spark.driver.cores=1 --conf spark.driver.memory=512M \ --conf spark.executor.instances=2 --conf spark.executor.cores=2 \ --conf spark.executor.memory=1G --name sparksql \ -f /tmp/5536_node.sql源码视角:参数如何被组装成提交命令
理解底层命令组装逻辑有助于排查提交类问题。核心方法populateSparkOptions()位于 SparkTask.java,其组装规则如下:
1. Master URL 的推导优先级
String masterUrl = StringUtils.isNotEmpty(sparkParameters.getMaster()) ? sparkParameters.getMaster() : onLocal ? deployMode : onNativeKubernetes ? SPARK_ON_K8S_MASTER_PREFIX + Config.fromKubeconfig(...).getMasterUrl() : SparkConstants.SPARK_ON_YARN;- 显式配置了
Master参数时,直接使用该值(如spark://localhost:7077); - 未配置时,
local模式用local作为 master; - 未配置但填写了 K8s 命名空间时,master 自动推导为
k8s://<kubeconfig 中的集群地址>; - 其余情况默认使用
yarn。
2. 部署方式
--deploy-mode仅在非local模式下追加;local模式不传该参数(见 SparkConstants.java 与测试用例中的断言)。
3. 资源参数一律以--conf形式注入
Driver / Executor 的资源配置在源码中均转换为--conf选项(SparkConstants.java):
--conf spark.driver.cores=%d--conf spark.driver.memory=%s--conf spark.executor.instances=%d--conf spark.executor.cores=%d--conf spark.executor.memory=%s
数值参数(核心数、Executor 数)仅在大于 0 时追加;内存参数仅在非空时追加。
4. YARN 队列
非local模式且可选参数中未包含--queue时,若配置了 YARN 队列则追加--queue <queue>(SparkConstants.java)。
5. Main Class 与主程序包
Java / Scala 类型且mainClass非空时追加--class <全限定名>;Python / SQL 类型不会追加该选项。非 SQL 类型会将主程序包在 Worker 本地的绝对路径追加为命令末尾的提交目标。
6. Spark on Kubernetes
当配置了命名空间时,插件会额外追加--conf spark.kubernetes.driver.label.<label>=<taskAppId>与--conf spark.kubernetes.namespace=<namespace>,用于驱动 Pod 打标签与指定命名空间(SparkConstants.java)。
7. 可选参数原样透传
可选参数字段的内容(如--jars、--files、--archives、--conf)会被原样拼接到命令中,是扩展 Spark 提交能力的最灵活入口。
此外,任务保存时的参数校验(checkParameters())与资源清单收集(getResourceFilesList()会把主程序包并入资源列表一并下载)逻辑可参考 SparkParameters.java,相关行为由 SparkParametersTest.java 验证。
注意事项
- Java 与 Scala 仅为类型标识:选择 Java 或 Scala 对执行没有任何区别,二者都通过
spark-submit提交 JAR 包; - Python 程序可忽略主函数的类:如果应用由 Python 开发,表单中的
主函数的类参数可以不填(源码中 Python 类型不会追加--class); - SQL 脚本参数仅 SQL 类型使用:Java / Scala / Python 类型可忽略
SQL 脚本参数; - SQL 不支持 cluster 模式:Spark SQL 任务目前只支持
client和local两种部署方式; - 环境依赖:Spark 任务最终依赖 Worker 机器上的
SPARK_HOME、JAVA_HOME以及提交目标集群(YARN / K8s / Standalone)的连接配置,务必在 dolphinscheduler_env.sh 中预先配置完整; - 资源文件可访问性:主程序包、SQL 脚本等资源必须位于资源中心且对执行租户有读取权限,执行时插件会将它们下载到 Worker 本地后再引用。
【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/GitHub_Trending/dol/dolphinscheduler
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考