拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Apache DolphinScheduler Amazon EMR Serverless 任务类型完整指南:配置、源码原理与实战示例

Apache DolphinScheduler Amazon EMR Serverless 任务类型完整指南配置、源码原理与实战示例【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/GitHub_Trending/dol/dolphinscheduler导读本文全面讲解 Apache DolphinScheduler 中的Amazon EMR ServerlessEMR_SERVERLESS任务类型从 DAG 画布上的任务创建、四个核心参数的配置到 Spark / Hive 两类作业的StartJobRunRequestJSON 实战示例、AWS 认证配置、作业状态机与故障转移机制。文中结合仓库源码EmrServerlessTask.java 与配套测试逐层剖析提交、轮询、取消与恢复的底层实现读完即可在真实集群中配置并排障 EMR Serverless 任务节点。Overview任务类型是什么Amazon EMR Serverless 任务类型用于向 Amazon EMR Serverless 应用提交并监控作业运行Job Run。与传统 EMR on EC2 不同EMR Serverless无需管理集群基础设施计算资源按需自动伸缩非常适合 Spark 与 Hive 负载。该任务类型在底层使用aws-java-sdkAWS SDK for Java其完整调用链为将表单中的 JSON 参数反序列化为 SDK 的StartJobRunRequest对象调用 StartJobRun API 向 AWS 提交作业获得jobRunId通过 GetJobRun API 轮询作业状态直到作业终止。从源码结构看该功能作为独立任务插件实现于dolphinscheduler-task-plugin/dolphinscheduler-task-emr-serverless模块插件名TaskChannelFactory.getName()返回值为EMR_SERVERLESS见 EmrServerlessTaskChannelFactory.java。它继承自AbstractRemoteTask远程任务基类按“提交应用 → 跟踪状态 → 映射退出码 → 支持取消/故障转移”的标准远程任务生命周期运行。创建任务进入Project Management - Project Name - Workflow Definition点击Create Workflow按钮进入 DAG 编辑页面。从工具栏将AmazonEMRServerless任务拖拽到画布上即完成节点创建。创建完成后在右侧Current node settings表单中配置下述参数。任务参数除默认参数外EMR Serverless 任务还需配置以下专属参数。通用默认参数节点名称、运行标志、失败重试次数、超时告警、延时执行时间、前置任务等请参阅 DolphinScheduler Task Parameters Appendix 的Default Task Parameters一节。参数说明Application IdEMR Serverless 应用 ID例如00fkht2eodujab09可从 EMR Serverless Console 获取Execution Role Arn作业执行所用 IAM 角色的 ARN例如arn:aws:iam::123456789012:role/EMRServerlessRole该角色需具备访问 S3、Glue 等服务的权限Job Name作业名称可选用于在 EMR Serverless 控制台中标识作业StartJobRunRequest JSON对应 StartJobRunRequest 中JobDriver与ConfigurationOverrides部分的 JSON示例见下文。注意JSON 中不要包含ApplicationId与ExecutionRoleArn它们会由上方表单参数自动注入参数校验规则源码级在 EmrServerlessParameters.java 中checkParameters()明确规定三个必填项缺一不可return StringUtils.isNotEmpty(applicationId) StringUtils.isNotEmpty(executionRoleArn) StringUtils.isNotEmpty(startJobRunRequestJson);对应地单元测试 EmrServerlessTaskTest.java 的testParametersCheck覆盖了“全空失败 / 仅有 applicationId 失败 / 缺 JSON 失败 / 三者齐备通过”四种场景。若参数缺失init()会抛出EmrServerlessTaskException(EMR Serverless task params are not valid)。JSON 的解析与字段注入原理buildStartJobRunRequest()是提交前最关键的一步见 EmrServerlessTask.java其处理顺序为参数占位符替换先调用ParameterUtils.convertParameterPlaceholders将 JSON 中的${...}占位符替换为工作流/任务的实际参数值如日期变量、自定义参数反序列化使用配置为UpperCamelCaseStrategy命名策略的 JacksonObjectMapper将 JSON 解析为StartJobRunRequest对象该策略使 JSON 中的JobDriver、ConfigurationOverrides等驼峰/大驼峰字段与 SDK 模型正确映射FAIL_ON_UNKNOWN_PROPERTIESfalse保证未知字段不导致解析失败强制注入用表单顶层的applicationId、executionRoleArn覆盖请求对象中同名字段这是文档强调“JSON 中不要写这两个字段”的源码依据命名回退若Job Name为空则自动使用 DolphinScheduler 的任务名称taskExecutionContext.getTaskName()幂等令牌自动设置clientToken taskInstanceId - System.currentTimeMillis()用于保证重复提交的幂等性。任务示例提交 Spark 作业以下示例演示创建EMR_SERVERLESS任务节点向 EMR Serverless 应用提交一个 Spark 作业。EntryPoint指向存放于 S3 的 JAREntryPointArguments传入输入/输出路径SparkSubmitParameters配置执行器资源{ JobDriver: { SparkSubmit: { EntryPoint: s3://my-bucket/scripts/my-spark-job.jar, EntryPointArguments: [ s3://my-bucket/input/, s3://my-bucket/output/ ], SparkSubmitParameters: --class com.example.MySparkApp --conf spark.executor.cores4 --conf spark.executor.memory8g --conf spark.executor.instances10 } }, ConfigurationOverrides: { MonitoringConfiguration: { S3MonitoringConfiguration: { LogUri: s3://my-bucket/emr-serverless-logs/ } } } }提交 Hive 作业以下示例演示创建EMR_SERVERLESS任务节点提交 Hive 查询。Query指向存放于 S3 的 SQL 脚本Parameters传入--hiveconf参数ApplicationConfiguration中的hive-site分类将 Hive 元数据指向 AWS Glue Data Catalog{ JobDriver: { HiveSQL: { Query: s3://my-bucket/scripts/my-hive-query.sql, Parameters: --hiveconf hive.exec.dynamic.partitiontrue --hiveconf hive.exec.dynamic.partition.modenonstrict } }, ConfigurationOverrides: { MonitoringConfiguration: { S3MonitoringConfiguration: { LogUri: s3://my-bucket/emr-serverless-logs/ } }, ApplicationConfiguration: [ { Classification: hive-site, Properties: { hive.metastore.client.factory.class: com.amazonaws.glue.catalog.metastore.AWSGlueDataCatalogHiveClientFactory } } ] } }在 JSON 中使用 DolphinScheduler 变量由于源码在反序列化前会先做参数占位符替换因此上述 JSON 可以嵌入 DolphinScheduler 的全局参数或自定义参数例如{ JobDriver: { SparkSubmit: { EntryPoint: s3://my-bucket/scripts/${spark_jar}.jar, EntryPointArguments: [ s3://my-bucket/input/${dt}/, s3://my-bucket/output/${dt}/ ] } } }这使得同一任务节点可在不同调度时间/不同业务流程中复用。AWS 认证配置EMR Serverless 任务从 DolphinScheduler 的aws.yaml配置文件读取 AWS 凭据读取段为conf/aws.yaml下的aws.emr部分。仓库内可参考的模板位于 aws.yaml其aws.emr段支持以下键键说明credentials.provider.type凭据提供方类型支持AWSStaticCredentialsProviderAK/SK与InstanceProfileCredentialsProviderIAM 角色access.key.id使用静态凭据时的 Access Key IDaccess.key.secret使用静态凭据时的 Secret Access KeyregionAWS 区域如us-east-1、cn-north-1使用 IAM 角色推荐当 DolphinScheduler Worker 节点运行在挂载了 IAM 角色的 EC2 实例上时可采用实例配置文件凭据aws: emr: credentials.provider.type: InstanceProfileCredentialsProvider region: us-east-1使用 Access Key如需使用 AK/SK 认证aws: emr: credentials.provider.type: AWSStaticCredentialsProvider access.key.id: your-access-key-id access.key.secret: your-secret-access-key region: us-east-1注意aws.emr段配置由 EMR on EC2 与 EMR Serverless 两种任务类型共享。客户端构建的源码细节在createEmrServerlessClient()EmrServerlessTask.java中凭据解析顺序为读取前缀为aws.emr.的全部配置项交给 AWSCredentialsProviderFactor.java 按credentials.provider.type创建AWSStaticCredentialsProvider或InstanceProfileCredentialsProvider若aws.emr.*配置缺失或无效回退到 AWS SDK 默认凭据链DefaultAWSCredentialsProviderChain即环境变量、系统属性、~/.aws/credentials 等标准 AWS 凭据发现机制支持通过配置项emr.serverless.endpoint或环境变量EMR_SERVERLESS_ENDPOINT自定义 endpoint便于使用 LocalStack 等 AWS 模拟服务进行本地测试——设置自定义 endpoint 且未指定 region 时默认使用us-east-1。作业状态转换EMR Serverless 作业提交后DolphinScheduler 每10 秒轮询一次作业状态SUBMITTED → PENDING → SCHEDULED → RUNNING → SUCCESS → FAILED → CANCELLED作业进入SUCCESS状态时任务标记为成功作业进入FAILED或CANCELLED状态时任务标记为失败若 DolphinScheduler 任务被 kill会自动调用 CancelJobRun API 取消正在运行的作业。状态轮询与退出码映射源码级源码 EmrServerlessTask.java 定义了“仍在进行中”的状态集合private static final HashSetString WAITING_STATES Sets.newHashSet( JobRunState.SUBMITTED.toString(), JobRunState.PENDING.toString(), JobRunState.SCHEDULED.toString(), JobRunState.RUNNING.toString());trackApplicationStatus()循环调用getJobRun()只要当前状态属于WAITING_STATES就TimeUnit.SECONDS.sleep(10)后继续轮询。最终状态到 DolphinScheduler 退出码的映射由mapStateToExitCode()完成EMR Serverless 最终状态DolphinScheduler 退出码SUCCESS成功EXIT_CODE_SUCCESSCANCELLED被杀EXIT_CODE_KILLFAILED及其他未知/空状态失败EXIT_CODE_FAILURE这一行为在单元测试 EmrServerlessTaskTest.java 中得到了完整验证testHandleSuccessSUBMITTED → RUNNING → SUCCESS、testHandleFailed、testHandleCancelled、testHandleFullLifecycle覆盖SUBMITTED → PENDING → SCHEDULED → RUNNING → SUCCESS全链路、testHandle_PollingFailure轮询中网络异常时任务置为失败。此外testSubmitError验证了 StartJobRun 抛错时任务抛出TaskException。任务取消cancelApplication()使用表单中的applicationId与已记录的jobRunId构造CancelJobRunRequest并调用cancelJobRun。若jobRunId为空尚未成功提交则直接跳过取消测试testCancelWithEmptyJobRunId验证了不会发起取消调用。故障转移Failover支持EMR Serverless 任务支持故障转移当 Worker 节点故障时新 Worker 可以通过appIds其中保存的是jobRunId恢复对运行中作业的追踪。具体机制如下submitApplication()在startJobRun成功后将返回的jobRunId通过setAppIds(jobRunId)持久化见 EmrServerlessTask.javatrackApplicationStatus()在jobRunId为空而appIds非空时从appIds恢复jobRunIdRecovered EMR Serverless jobRunId from appIds随后正常轮询状态。测试testFailoverRecovery完整模拟了这一场景提交作业后新建一个任务实例注入之前持久化的appIds验证恢复后不会再次调用startJobRun而是直接通过getJobRun追踪到已完成状态并返回成功。注意事项Application Id必须对应一个已存在的 EMR Serverless 应用通过 AWS Console 或 API 创建且处于STARTED或CREATED状态Execution Role至少需要以下最小权限emr-serverless:StartJobRun、emr-serverless:GetJobRun、emr-serverless:CancelJobRun以及作业运行所需的 S3、Glue 等数据访问权限StartJobRunRequest JSON中不应包含ApplicationId或ExecutionRoleArn字段——它们会从表单参数自动注入源码中request.setApplicationId(...)与request.setExecutionRoleArn(...)会直接覆盖 JSON 中的同名值三个必填参数applicationId、executionRoleArn、startJobRunRequestJson任一缺失任务初始化即失败若在本地开发或测试环境如 LocalStack中使用可通过emr.serverless.endpoint配置项或EMR_SERVERLESS_ENDPOINT环境变量指向模拟服务。相关源码与文档索引任务实现EmrServerlessTask.java参数模型EmrServerlessParameters.java插件注册EmrServerlessTaskChannelFactory.java 与 EmrServerlessTaskChannel.java单元测试覆盖全生命周期、取消、故障转移、参数校验EmrServerlessTaskTest.javaAWS 凭据配置模板aws.yaml凭据工厂实现AWSCredentialsProviderFactor.java通用任务参数DolphinScheduler Task Parameters Appendix【免费下载链接】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),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门