ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

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

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/dolphinschedulerSpark 任务节点是 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.coresN可按实际生产环境调整。Driver 内存大小 (Driver memory size)设置 Driver 内存大小对应--conf spark.driver.memorySIZE。Executor 数量 (Number of Executor)设置 Executor 个数对应--conf spark.executor.instancesN。Executor 内存大小 (Executor memory size)设置每个 Executor 的内存大小对应--conf spark.executor.memorySIZE。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-wordcountDriver / 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.cores1 --conf spark.driver.memory512M \ --conf spark.executor.instances2 --conf spark.executor.cores2 \ --conf spark.executor.memory1G --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.javaFILE 模式从资源中心读取.sql文件内容多个文件时默认取第一个源码会打印告警more than 1 files detected, use the first one by defaultSCRIPT 模式使用表单中的内联脚本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.cores1 --conf spark.driver.memory512M \ --conf spark.executor.instances2 --conf spark.executor.cores2 \ --conf spark.executor.memory1G --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 queueSparkConstants.java。5. Main Class 与主程序包Java / Scala 类型且mainClass非空时追加--class 全限定名Python / SQL 类型不会追加该选项。非 SQL 类型会将主程序包在 Worker 本地的绝对路径追加为命令末尾的提交目标。6. Spark on Kubernetes当配置了命名空间时插件会额外追加--conf spark.kubernetes.driver.label.labeltaskAppId与--conf spark.kubernetes.namespacenamespace用于驱动 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 类型不会追加--classSQL 脚本参数仅 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),仅供参考
返回列表