news 2026/9/15 19:58:15

Apache DolphinScheduler Spark 任务节点完全指南:spark-submit 与 Spark SQL 提交原理与实战

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache DolphinScheduler Spark 任务节点完全指南:spark-submit 与 Spark SQL 提交原理与实战

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 集群提交作业:

  1. spark-submit方式:适用于 Java / Scala / Python 类型的 Spark 应用,向集群提交主程序(JAR 包或 Python 文件)并附带启动参数。
  2. 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()方法负责选择命令:当ProgramTypeSQL时走spark-sql,否则走spark-submit,随后把populateSparkOptions()组装出的参数拼接为最终命令字符串。该逻辑与 ProgramType.java 中定义的JAVA / SCALA / PYTHON / SQL四种程序类型一一对应。

创建 Spark 任务

在 DolphinScheduler Web UI 中按以下步骤创建 Spark 任务节点:

  1. 进入项目管理 -> 项目名称 -> 工作流定义,点击创建工作流按钮进入 DAG 编辑页面;
  2. 从左侧工具栏拖拽SPARK节点到画布中;
  3. 双击节点,在右侧表单中按下文"任务参数详解"填写配置;
  4. 保存节点并完成 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,例如yarnspark://localhost:7077或 Kubernetes 集群地址。留空时由插件按规则推导(见"命令组装"一节)。
主程序包 (Main jar package)通过资源中心上传的 Spark 主程序 JAR 包(Java / Scala 使用)或 Python 文件(Python 类型使用)。
SQL 脚本 (SQL scripts)Spark SQL 执行的.sql文件(SQL 类型使用),从资源中心选择。
部署方式 (Deployment mode)spark-submit支持clusterclientlocal三种模式;spark-sql支持clientlocal两种模式。
命名空间(集群)(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 承载(mainJarmainClassmasterdeployModemainArgsdriverCoresdriverMemorynumExecutorsexecutorCoresexecutorMemoryappNameyarnQueueothersprogramTyperawScriptnamespaceresourceListsqlExecutionType等字段),任务保存时通过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_HOMESPARK_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_HOMEHADOOP_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
  • 部署方式clientcluster(提交到 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 支持clientlocal两种模式,不支持 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 任务目前只支持clientlocal两种部署方式;
  • 环境依赖:Spark 任务最终依赖 Worker 机器上的SPARK_HOMEJAVA_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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/15 19:55:48

MV3浏览器插件工程化:从架构演进到端侧AI实战指南

入行这么多年&#xff0c;我经常见到有人把浏览器插件想成一个小脚本&#xff1a;改改页面样式、往页面里塞一段逻辑&#xff0c;完事。可当你真把一个插件从 MVP 推到线上&#xff0c;被用户报了一堆跨标签页不同步、后台任务被回收、敏感数据泄漏的问题之后&#xff0c;会明白…

作者头像 李华
网站建设 2026/9/15 19:54:52

mise声明式环境管理:统一Java/Node/Maven多版本开发环境

1. 项目概述&#xff1a;当开发环境管理从“手动拼凑”走向“声明式交付”最近三个月&#xff0c;我彻底把本地开发环境的控制权交给了mise——不是简单地换了个工具&#xff0c;而是重构了整个工程化基础设施的认知逻辑。过去写 Java 项目要配 JDK、配 Maven、配JAVA_HOME&…

作者头像 李华
网站建设 2026/9/15 19:53:55

MATLAB实现SOFT立体视觉里程计:从特征跟踪到局部地图优化

简介&#xff1a;面向机器人技术、立体视觉与视觉里程计研究者的MATLAB实现&#xff0c;基于SOFT算法完成特征选择与跟踪&#xff0c;并估计相机运动轨迹。代码已在MATLAB R2018a上测试&#xff0c;依赖并行处理与计算机视觉工具箱&#xff0c;同时给出特征处理、匹配、选择及运…

作者头像 李华
网站建设 2026/9/15 19:52:33

金融数据入湖架构设计与实践指南

1. 金融数据入湖的背景与挑战金融行业正面临数据爆炸式增长的时代。根据国际数据公司&#xff08;IDC&#xff09;的统计&#xff0c;全球金融服务业数据量每年以40%以上的速度增长&#xff0c;而传统的数据仓库架构已经难以应对这种海量、多样化的数据处理需求。数据湖&#x…

作者头像 李华
网站建设 2026/9/15 19:52:26

衍射型多焦点人工晶体倾斜入射效应:原理、测试与仿真

前阵子在一套光学测试台上测衍射型多焦点人工晶体&#xff0c;样品转了十几度&#xff0c;像面上的焦点分布立刻就不一样了——焦点位置在跑&#xff0c;能量比在变&#xff0c;边缘光斑还拖出了尾影。这个现象如果没搞懂背后的光栅原理&#xff0c;很容易在测试和临床随访里被…

作者头像 李华
网站建设 2026/9/15 19:50:03

copilot.vim 实战指南:在 Vim/Neovim 中配置与使用 GitHub Copilot

copilot.vim 实战指南&#xff1a;在 Vim/Neovim 中配置与使用 GitHub Copilot 【免费下载链接】copilot.vim Neovim plugin for GitHub Copilot 项目地址: https://gitcode.com/GitHub_Trending/co/copilot.vim 本指南以 copilot.vim 插件的官方 README.md 为骨架&…

作者头像 李华