Apache DolphinScheduler 集成阿里云 EMR Serverless Spark 任务插件:从数据源配置到四种作业提交实战
任务调度大数据后端前端【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址https://gitcode.com/gh_mirrors/do/dolphinscheduler点击查看免费下载本指南面向在 Apache DolphinScheduler 中编排阿里云 EMR Serverless Spark 作业的开发者完整讲解ALIYUN_SERVERLESS_SPARK任务插件从数据源创建、任务节点配置到 JAR / SQL / PySpark 作业提交的全流程并深入源码层剖析作业提交、状态轮询与取消的底层实现帮助读者在无需自建 Spark 集群的前提下将 Serverless Spark 任务无缝编排进工作流。插件简介Aliyun EMR Serverless Spark是 Apache DolphinScheduler 提供的一种远程任务插件用于向阿里云 EMR Serverless Spark 服务提交 Spark 作业。它的核心价值在于不需要自行维护 Spark 集群只需配置好阿里云账号的访问凭证AccessKey与区域信息即可在 DolphinScheduler 的工作流中以节点形式直接提交 JAR、Python 或 SQL 类型的 Spark 作业作业运行在阿里云托管的 Serverless Spark 环境上。在源码层面该插件由两个模块协同工作数据源插件dolphinscheduler-datasource-aliyunserverlessspark负责数据源链接的创建、参数校验与连通性检测任务插件dolphinscheduler-task-aliyunserverlessspark负责将任务参数组装为阿里云 OpenAPI 请求并提交、跟踪、取消作业。一、创建数据源链接在提交任务之前需要先创建ALIYUN_SERVERLESS_SPARK类型的数据源DolphinScheduler 中称为“链接”用于统一管理阿里云账号凭证多个任务节点可以复用同一数据源。操作路径点击数据源 - 创建数据源 - ALIYUN_SERVERLESS_SPARK进入创建表单。在表单中需要填写以下参数参数说明Datasource Name数据源名称用于在工作流中标识该链接Access Key Id阿里云账号的 AccessKey ID用于调用 OpenAPIAccess Key Secret阿里云账号的 AccessKey Secret用于签名鉴权Region Id地域 ID例如cn-hangzhouServerless Spark 服务所在区域填写完成后点击确认保存数据源。从源码实现看数据源的参数模型定义在 AliyunServerlessSparkConnectionParam.java除上述四个字段外还支持可选的自定义endpoint。当 endpoint 为空时插件会按emr-serverless-spark.{regionId}.aliyuncs.com的模板自动拼接见 AliyunServerlessSparkConstants.java 中的ENDPOINT_TEMPLATE。数据源创建时会做参数合法性校验与连通性检查AliyunServerlessSparkDataSourceProcessor.java 中的checkDatasourceParam要求Region Id与Access Key Id均不能为空否则抛出IllegalArgumentExceptioncheckDataSourceConnectivity会通过阿里云 SDK 客户端实际发起一次连通性探测checkConnect失败则返回 false保证只有凭证有效、网络可达的数据源才能保存成功。二、创建任务节点数据源就绪后即可在工作流中创建任务节点点击项目 - 工作流定义 - 创建工作流将ALIYUN_SERVERLESS_SPARK任务从左侧任务面板拖拽到画板中。在节点配置表单中填写任务参数后点击确认完成节点创建之后即可像普通节点一样参与工作流调度与依赖编排。三、任务参数详解ALIYUN_SERVERLESS_SPARK节点的任务参数如下默认参数说明请参考 DolphinScheduler 任务参数附录 的“默认任务参数”一栏任务参数描述Datasource types链接类型应该选择ALIYUN_SERVERLESS_SPARKDatasource instancesALIYUN_SERVERLESS_SPARK链接实例workspace idAliyun Serverless Spark工作空间 IDresource queue idAliyun Serverless Spark任务队列 IDcode typeAliyun Serverless Spark任务类型可以是JAR、PYTHON或者SQLjob nameAliyun Serverless Spark任务名entry point任务代码JAR 包、PYTHON / SQL 脚本的位置支持 OSS 中的文件entry point arguments主程序入口参数spark submit parametersSpark-submit 相关参数engine release versionSpark 引擎版本is productionSpark 任务是否运行在生产环境中这些表单字段与 AliyunServerlessSparkParameters.java 中的字段一一对应workspaceId、resourceQueueId、codeType、jobName、engineReleaseVersion、entryPoint、entryPointArguments、sparkSubmitParameters、isProduction此外还有datasource数据源 ID与type链接类型。关键参数的填写要点engine release version引擎版本可直接填写官方发布版本例如esr-2.1-native (Spark 3.3.1, Scala 2.12, Native Runtime)。若留空插件会使用默认引擎版本见 AliyunServerlessSparkConstants.java 中的DEFAULT_ENGINEentry point arguments入口参数多个参数之间使用#作为分隔符插件在提交时按该分隔符拆分后传给 Spark 作业对应常量ENTRY_POINT_ARGUMENTS_DELIMITER。例如 SQL 类型任务中-e#show tables;show tables;会被拆分为-e与show tables;show tables;两个参数与 SparkSQLCLIDriver 的用法对应is production是否生产环境开启后作业会附带environmentproduction标签否则为environmentdev标签便于在阿里云控制台区分生产与开发作业对应 AliyunServerlessSparkTask.java 中的标签构建逻辑。四、源码视角任务提交、状态跟踪与取消理解底层实现有助于排查“任务提交成功但状态不更新”“如何区分运行环境”等实际问题。任务执行的核心类是 AliyunServerlessSparkTask.java它继承AbstractRemoteTask通过阿里云 EMR Serverless Spark OpenAPI SDKcom.aliyun.emr_serverless_spark20230808完成全部交互。1. 初始化init任务初始化时依次完成从TaskExecutionContext解析任务参数 JSON 为AliyunServerlessSparkParameters解析失败则抛出AliyunServerlessSparkTaskException通过datasource字段从资源参数中取出对应的数据源构建AliyunServerlessSparkConnectionParam读取accessKeyId、accessKeySecret、regionId、endpoint调用buildAliyunServerlessSparkClient创建阿里云 SDK 客户端endpoint 为空时按emr-serverless-spark.{regionId}.aliyuncs.com拼接并用 AccessKey 构造Config将初始运行状态置为RunState.Submitted。2. 提交与轮询handlehandle方法是任务的主流程分为两步提交作业组装StartJobRunRequest设置regionId、resourceQueueId、codeType、name即 job name、releaseVersion引擎版本并写入两个标签——environment值为production或dev由isProduction决定与workflowtrue标识该作业由工作流提交。随后通过startJobRunWithOptions提交返回的jobRunId会写入 appIds用于在 DolphinScheduler 界面上关联与追踪轮询状态进入循环每10 秒调用一次getJobRun查询作业状态直到状态进入终态。状态枚举定义在 RunState.javaSubmitted、Pending、Running、Success、Failed、Cancelling、Cancelled、CancelFailed其中Success、Failed、Cancelled为终态isFinal。轮询结束后终态会被映射为 DolphinScheduler 的退出码mapFinalStateToExitCodeSuccess-EXIT_CODE_SUCCESS任务成功Failed-EXIT_CODE_KILL其他异常状态 -EXIT_CODE_FAILURE。单元测试 AliyunServerlessSparkTaskTest.java 覆盖了上述流程mock 客户端返回Success状态后断言最终退出码为EXIT_CODE_SUCCESS同时也验证了任务参数 JSON 的结构workspaceId、resourceQueueId、codeType、entryPoint、sparkSubmitParameters等字段。3. 取消作业cancelApplication当工作流被停止、超时或用户主动 kill 任务时插件会调用cancelJobRun接口按jobRunId取消远程作业对应 OpenAPI 的CancelJobRunRequest。4. 连通性测试在数据源模块中AliyunServerlessSparkDataSourceProcessor.java 的checkDataSourceConnectivity通过AliyunServerlessSparkClientWrapper.checkConnect验证凭证有效性这是创建数据源时“测试连接”按钮背后的实现。五、四种作业提交示例以下示例均基于cn-hangzhou区域AccessKey 请替换为实际值示例代码与 OSS 路径取自插件单元测试中的典型配置。示例一提交 JAR 类型任务以官方 Spark 示例包spark-examples_2.12-3.3.1.jar运行 SparkPi 为例参数名参数值 / 按钮操作region idcn-hangzhouaccess key idyour-access-key-idaccess key secretyour-access-key-secretresource queue idroot_queuecode typeJARjob nameds-emr-spark-jarentry pointoss://datadev-oss-hdfs-test/spark-resource/examples/jars/spark-examples_2.12-3.3.1.jarentry point arguments100spark submit parameters--class org.apache.spark.examples.SparkPi --conf spark.executor.cores4 --conf spark.executor.memory20g --conf spark.driver.cores4 --conf spark.driver.memory8g --conf spark.executor.instances1engine release versionesr-2.1-native (Spark 3.3.1, Scala 2.12, Native Runtime)is production请您将按钮打开要点JAR 任务必须通过spark submit parameters指定--class主类entry point指向 OSS 上的 JAR 包entry point arguments中的100是传给主程序的参数SparkPi 的迭代次数。示例二提交 SQL 类型任务以 SparkSQLCLIDriver 直接执行 SQL 为例参数名参数值 / 按钮操作region idcn-hangzhouaccess key idyour-access-key-idaccess key secretyour-access-key-secretresource queue idroot_queuecode typeSQLjob nameds-emr-spark-sql-1entry point任意非空值entry point arguments-e#show tables;show tables;spark submit parameters--class org.apache.spark.sql.hive.thriftserver.SparkSQLCLIDriver --conf spark.executor.cores4 --conf spark.executor.memory20g --conf spark.driver.cores4 --conf spark.driver.memory8g --conf spark.executor.instances1engine release versionesr-2.1-native (Spark 3.3.1, Scala 2.12, Native Runtime)is production请您将按钮打开要点entry point需要填入任意非空值占位即可主类必须为SparkSQLCLIDriverentry point arguments中-e表示直接执行内联 SQL#是参数分隔符因此-e#show tables;show tables;实际等价于-e show tables;show tables;。示例三提交 OSS 中的 SQL 脚本任务与示例二的区别在于通过-f指定 OSS 上的 SQL 脚本文件参数名参数值 / 按钮操作region idcn-hangzhouaccess key idyour-access-key-idaccess key secretyour-access-key-secretresource queue idroot_queuecode typeSQLjob nameds-emr-spark-sql-2entry point任意非空值entry point arguments-f#oss://datadev-oss-hdfs-test/spark-resource/examples/sql/show_db.sqlspark submit parameters--class org.apache.spark.sql.hive.thriftserver.SparkSQLCLIDriver --conf spark.executor.cores4 --conf spark.executor.memory20g --conf spark.driver.cores4 --conf spark.driver.memory8g --conf spark.executor.instances1engine release versionesr-2.1-native (Spark 3.3.1, Scala 2.12, Native Runtime)is production请您将按钮打开要点-f#oss://.../show_db.sql会被拆分为-f与 OSS 脚本路径SparkSQLCLIDriver 的-f选项表示从文件读取 SQL从而支持把脚本托管在 OSS 上统一管理。示例四提交 PySpark 任务参数名参数值 / 按钮操作region idcn-hangzhouaccess key idyour-access-key-idaccess key secretyour-access-key-secretresource queue idroot_queuecode typePYTHONjob nameds-emr-spark-pythonentry pointoss://datadev-oss-hdfs-test/spark-resource/examples/src/main/python/pi.pyentry point arguments100spark submit parameters--conf spark.executor.cores4 --conf spark.executor.memory20g --conf spark.driver.cores4 --conf spark.driver.memory8g --conf spark.executor.instances1engine release versionesr-2.1-native (Spark 3.3.1, Scala 2.12, Native Runtime)is production请您将按钮打开要点entry point指向 OSS 上的.py脚本与 JAR 类型不同PySpark 任务无需在 spark submit parameters 中指定--class。六、常见问题与注意事项凭证错误或区域不匹配任务在init阶段构建客户端失败会直接抛出异常请核对数据源中的 AccessKey 与 Region Id 是否一致若使用 VPC 内网可在数据源中显式配置endpoint如emr-serverless-spark-vpc.cn-hangzhou.aliyuncs.com否则默认使用公网 endpoint。任务长时间处于 Submitted / Pending插件每 10 秒轮询一次远程状态若集群资源排队作业会长时间停留在Pending属正常现象可通过阿里云控制台结合workflowtrue与environment标签快速定位该作业。SQL 任务入口参数分隔符多个参数务必用#分隔如-e#show tables;不要使用空格否则整个字符串会被当作单个参数传入。is production 开关开启后作业携带environmentproduction标签建议生产工作流开启方便与开发环境作业区分与审计。任务失败排查任务失败退出码为EXIT_CODE_KILLFailed状态可在 DolphinScheduler 任务实例日志中查看错误堆栈远程作业的具体失败原因以阿里云 EMR Serverless Spark 控制台的作业运行详情为准。延伸阅读任务通用参数超时、重试、资源等默认参数说明DolphinScheduler 任务参数附录数据源插件实现dolphinscheduler-datasource-aliyunserverlessspark任务插件实现dolphinscheduler-task-aliyunserverlessspark插件单元测试AliyunServerlessSparkTaskTest.java赞分享任务调度大数据后端前端【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址https://gitcode.com/gh_mirrors/do/dolphinscheduler点击查看免费下载相关推荐Apache DolphinScheduler 集成阿里云 EMR Serverless Spark 任务插件从数据源配置到 JAR / SQL / PySpark 作业提交实战Apache DolphinScheduler 集成阿里云 EMR Serverless Spark 任务插件从数据源配置到 JAR / SQL / PySp任务调度数据编排工作流自动化后端大数据CANN/.gitcode镜像修订工具revise img 根据 CI 仓库默认 cann/.gitcode 中 image conf/target_branch /images.yaml 的任务调度数据编排工作流自动化后端大数据在 Apache DolphinScheduler 中编排 Aliyun EMR Serverless Spark 任务数据源配置、任务参数与源码级运行原理在 Apache DolphinScheduler 中编排 Aliyun EMR Serverless Spark 任务数据源配置、任务参数与源码级运行原理任务调度大数据后端前端创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考