AWS Step Functions结合EMR Serverless的作业编排实现方法
EMR Serverless 提供与传统EMR Step能力完全等价的AWS Step Functions官方原生集成,无需依赖Airflow等第三方编排框架,也不需要自行编写Lambda做作业状态轮询,原生支持同步等待作业执行完成的模式,可以直接替换原有arn:aws:states:::elasticmapreduce:addStep.sync类型的任务节点。
原有传统EMR同步提交作业的资源ARN为arn:aws:states:::elasticmapreduce:addStep.sync,对应EMR Serverless的同步作业提交资源ARN为:
arn:aws:states:::emr-serverless:startJobRun.sync
该集成由AWS官方维护,运行逻辑和传统EMR同步集成完全一致:Step Functions会自动轮询EMR Serverless作业的运行状态,作业成功执行完成后才会流转到下一个流程节点,作业执行失败、超时会直接抛出对应错误,完全匹配现有同步执行的流程要求。
原有addStep.sync节点强制要求传入ClusterId参数,EMR Serverless场景不存在临时集群概念,不需要传入集群ID,替换为以下必填参数即可:
applicationId:提前创建好的EMR Serverless应用ID,对应固定版本的Spark运行环境,替代原有流程中临时创建的EMR集群executionRoleArn:EMR Serverless作业运行时使用的IAM角色ARN,需要配置作业访问S3、Glue等关联服务的对应权限jobDriver:作业驱动配置,直接把原有通过spark-submit提交的Scala作业参数(主类路径、Jar包S3地址、Spark提交参数、业务程序入参)填入对应的Spark提交配置块即可- 可选参数:作业运行的资源规格(Driver、Executor的vCPU、内存配置)、日志输出路径、自定义Spark配置项等
原来的四步流程可以直接省去创建EMR集群、终止EMR集群两个环节,调整后的流程为:
- 运行原有前置Lambda函数,完成数据校验、输出路径清理等前置准备
- 调用
emr-serverless:startJobRun.sync节点提交Spark作业,同步等待作业执行完成
注意:EMR Serverless应用是长期留存的资源,不需要每次作业前创建、作业后销毁,提前创建一次对应Spark版本的应用即可。作业提交后服务会自动按需分配计算资源,作业结束后自动释放计算资源,不存在空闲集群的运行成本。
要保证Step Functions正常调用EMR Serverless提交作业,需要给Step Functions的执行角色附加对应权限:
- 允许对目标EMR Serverless应用执行
emr-serverless:StartJobRun操作 - 允许对目标EMR Serverless应用执行
emr-serverless:GetJobRun操作(同步模式下Step Functions需要通过该接口轮询作业状态) - EMR Serverless作业关联的执行角色,权限配置和原有EMR集群上提交作业的EC2实例角色保持一致即可,不需要做额外调整。
内容的提问来源于stack exchange,提问作者smishra

