如何自动化向运行中的AWS EMR集群添加pyspark-submit步骤?
AWS EMR长期运行集群PySpark脚本自动化提交方案
基于你给出的约束条件,以下两种方案均可直接实现本地PyCharm端提交PySpark任务到目标集群,完全符合权限要求。
方案1:boto3调用EMR原生API提交(最灵活,支持全自定义逻辑)
该方案直接复用EMR控制台Steps添加的底层接口,和团队手动在控制台添加Step的权限逻辑完全一致,无需额外资源。
前置准备
- 本地PyCharm环境安装boto3:
pip install boto3 - 配置本地AWS凭证(可通过
~/.aws/credentials文件、环境变量或AWS SSO完成配置),确保凭证拥有emr:AddJobFlowSteps、s3:PutObject、s3:GetObject权限。
实现步骤
- 把PySpark代码上传到你有权限的S3路径,支持本地文件上传或直接上传代码片段
- 调用EMR的AddJobFlowSteps接口提交Spark任务
完整代码示例
import boto3 from io import BytesIO # 初始化AWS客户端,替换为你的集群所属区域 region = 'cn-north-1' s3_client = boto3.client('s3', region_name=region) emr_client = boto3.client('emr', region_name=region) # 自定义配置项,替换为实际值 CLUSTER_ID = 'j-XXXXXXXX' # 目标EMR集群ID S3_BUCKET = 'your-bucket-name' # 有权限的S3桶名 S3_SCRIPT_KEY = 'pyspark_jobs/test_job.py' # S3上存储脚本的路径 # --- 可选1:上传本地PySpark脚本到S3 --- # local_script_path = './my_pyspark_script.py' # s3_client.upload_file(local_script_path, S3_BUCKET, S3_SCRIPT_KEY) # --- 可选2:直接上传PySpark代码片段到S3,无需本地存文件 --- pyspark_code = """ from pyspark.sql import SparkSession if __name__ == "__main__": spark = SparkSession.builder.appName("AutoSubmitTest").getOrCreate() # 你的业务逻辑 df = spark.range(1000) print(f"数据量:{df.count()}") spark.stop() """ s3_client.upload_fileobj( Fileobj=BytesIO(pyspark_code.encode('utf-8')), Bucket=S3_BUCKET, Key=S3_SCRIPT_KEY ) # 提交EMR Step submit_resp = emr_client.add_job_flow_steps( JobFlowId=CLUSTER_ID, Steps=[{ 'Name': 'Auto_Submit_PySpark_Job', 'ActionOnFailure': 'CONTINUE', # 任务失败不影响集群运行,禁止改为TERMINATE_CLUSTER 'HadoopJarStep': { 'Jar': 'command-runner.jar', # EMR内置执行器,无需额外上传 'Args': [ 'spark-submit', '--deploy-mode', 'cluster', # 可根据需要添加其他spark配置,如--executor-memory、--num-executors等 f's3://{S3_BUCKET}/{S3_SCRIPT_KEY}', # 此处可添加PySpark脚本需要的入参,按顺序传入即可 ] } }] ) # 输出提交的Step ID,可用于后续查询运行状态 step_id = submit_resp['StepIds'][0] print(f"任务提交成功,Step ID:{step_id}")
任务状态查询
你可以调用emr_client.list_steps(ClusterId=CLUSTER_ID, StepIds=[step_id])接口查询运行状态,也可以直接在EMR控制台Steps tab查看运行日志,日志默认存储在集群配置的S3日志路径下,你有S3权限可直接访问。
方案2:PyCharm专业版AWS Toolkit可视化提交
如果你不想自己编写boto3代码,可直接使用PyCharm专业版支持的AWS Toolkit插件实现可视化提交:
- 第一步:在PyCharm插件市场搜索安装AWS Toolkit,完成后配置你的AWS凭证和对应区域
- 第二步:在插件的EMR列表中找到目标集群,右键选择「Submit Spark job」,按提示选择本地PySpark脚本、配置spark参数即可提交,插件会自动完成脚本上传S3、提交Step的全流程。
注意事项
- 确保EMR集群关联的EC2实例角色拥有你脚本中用到的所有资源访问权限,例如读写S3业务数据、访问其他AWS服务的权限
- 提交Step时
ActionOnFailure参数不要设置为TERMINATE_CLUSTER,避免任务失败误关闭长期运行的集群 - 如果你的PySpark任务有依赖包,可将依赖打包上传到S3,在spark-submit参数中通过
--py-files指定依赖包的S3路径即可
内容的提问来源于stack exchange,提问作者RandyB
相关产品推荐
相关产品推荐

