You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.29 00:48:02