如何通过Airflow连接AWS EMR集群上运行的Notebook
解决方案
核心实现逻辑
EMR没有提供直接运行预设Notebook的原生Step类型,我们可以通过jupyter nbconvert命令在EMR Step中执行ipynb文件,完全兼容你现有的EmrAddStepsOperator调用逻辑,无需额外引入新的Operator。
前置检查
- 预设的EMR Notebook已导出为
.ipynb格式并上传到S3桶的可访问路径 - EMR集群实例角色(EC2 Instance Profile)已配置S3读写权限,可读取Notebook文件、写入运行结果
- EMR集群已安装
nbconvert依赖(EMR 6.x以上版本默认预装,未预装可补充安装步骤) - Airflow绑定的
aws_default连接角色具备EMR Step提交权限
代码修改说明
1. 补全缺失的依赖导入
原代码遗漏了时间相关依赖的导入,在代码开头新增如下导入语句:
from datetime import datetime, timedelta
2. 新增Notebook相关配置项
在原有配置区补充Notebook对应的路径配置:
# 新增Notebook相关配置 s3_notebook = "notebooks/your_predefined_notebook.ipynb" # 替换为你的预设Notebook的S3路径 s3_notebook_output = "notebook_output/" # 替换为Notebook运行结果的S3输出路径
3. 新增执行Notebook的Step到SPARK_STEPS
如果你的集群已预装nbconvert,直接在原有SPARK_STEPS数组末尾添加如下步骤即可,参数化写法和你现有代码风格完全统一:
{ "Name": "Run predefined EMR Notebook", "ActionOnFailure": "CANCEL_AND_WAIT", "HadoopJarStep": { "Jar": "command-runner.jar", "Args": [ "bash", "-c", """ # 从S3拉取Notebook到集群本地 aws s3 cp s3://{{ params.BUCKET_NAME }}/{{ params.s3_notebook }} /tmp/ # 执行Notebook并输出html格式的运行结果 jupyter nbconvert --execute --to html /tmp/$(basename {{ params.s3_notebook }}) --output /tmp/notebook_output.html # 运行结果回传到S3 aws s3 cp /tmp/notebook_output.html s3://{{ params.BUCKET_NAME }}/{{ params.s3_notebook_output }}/ """ ] } }
如果集群未预装nbconvert,在上述步骤之前额外加一个安装依赖的Step即可:
{ "Name": "Install nbconvert dependency", "ActionOnFailure": "CANCEL_AND_WAIT", "HadoopJarStep": { "Jar": "command-runner.jar", "Args": ["pip", "install", "nbconvert"] } }
4. 补充参数到step_adder的params配置
在现有step_adder的params字段中新增Notebook相关的参数键:
params={ "BUCKET_NAME": BUCKET_NAME, "s3_data": s3_data, "s3_script": s3_script, "s3_clean": s3_clean, "s3_notebook": s3_notebook, # 新增 "s3_notebook_output": s3_notebook_output # 新增 },
异常排查
如果运行时报权限错误,按如下规则检查即可:
- 确认Airflow使用的
aws_default连接的AK/SK或角色有emr:AddSteps、emr:DescribeStep权限 - 确认EMR集群的EC2实例角色有对应S3桶的
s3:GetObject、s3:PutObject权限 - 确认EMR集群的安全组允许出网访问S3端点
内容的提问来源于stack exchange,提问作者Mobeen
相关产品推荐
相关产品推荐

