Apache Airflow AWS Batch Executor 实战指南:以 Amazon Batch 弹性运行工作流
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
Apache Airflow 的 AWS Batch Executor(AwsBatchExecutor)将 Scheduler 派发的每一个任务封装为一次 AWS Batch 作业,在由 Amazon Batch 动态供给的独立容器中执行,从而把任务调度与资源供给彻底解耦。本指南围绕该 Executor 的配置项、容器镜像构建、IAM 权限、远程日志与端到端快速上手指南展开,结合 amazon provider 源码说明其底层调用链,读完即可在自己的 Airflow 环境中落地一套可弹性伸缩、按需付费的批处理执行体系。
AWS Batch Executor 是什么
AWS Batch Executor 是 Apache Airflow 的一种 Executor 实现,位于 amazon provider 包中,核心类为AwsBatchExecutor(batch_executor.py)。它的工作模式是:Airflow Scheduler 为每个任务生成执行命令,Executor 将该命令作为containerOverrides.command提交给 AWS Batch 的submit_jobAPI,之后周期性地通过describe_jobs按 job-id 轮询作业状态,并在作业失败或成功时回调 Airflow 更新任务实例状态。
相比在同一台主机上并发运行任务的本地 Executor,将每个任务交给 AWS Batch 运行带来以下收益:
- 可扩展性与更低成本:AWS Batch 能够按需动态供给执行任务所需的资源,并根据负载自动扩缩容,资源利用率高、总体成本更低。
- 任务队列与优先级:AWS Batch 提供 Job Queue 概念,可对任务执行进行优先级编排。当多个任务同时被调度时,能按期望的优先级顺序执行。
- 灵活性:AWS Batch 支持 Fargate(ECS)、EC2 和 EKS 三种计算环境,配合对计算环境资源(vCPU、内存、GPU 等)的精细定义,用户可为工作负载选择最合适的执行环境。
- 快速任务执行:通过维护处于活跃状态的 worker,提交给 Batch 的任务能够迅速被执行;就绪的 worker 几乎无启动延迟,特别适合对时效敏感或需要近实时处理的工作负载。
从源码看执行生命周期
Executor 的核心循环在sync()方法中(batch_executor.py),由 Scheduler 的心跳周期性地触发,包含两个阶段:
sync_running_jobs():将当前所有活动作业的 job-id 以每批最多 99 个(DESCRIBE_JOBS_BATCH_SIZE,受 AWSdescribe_jobs上限约束)调用describe_jobs查询状态;attempt_submit_jobs():把等待队列中的作业逐个提交到 Batch。
Batch 作业状态与 Airflow 状态的映射关系定义在 utils.py:SUBMITTED/PENDING/RUNNABLE/STARTING映射为 Airflow 的QUEUED,RUNNING映射为RUNNING,SUCCEEDED映射为SUCCESS,FAILED映射为FAILED。submit_job与describe_jobs的响应分别由 boto_schema.py 中的 marshmallow Schema 校验并转成领域对象。
值得注意的是,AWS Batch 标记为FAILED的作业并不等同于 Airflow 层面的任务失败:如果容器能正常启动并运行 Airflow 进程,之后的 DAG 级失败由 Airflow 自行捕获处理;而容器启动之前发生的失败(如 Batch API 故障、容器配置错误)才会被 Batch 标记为FAILED,此时 Executor 会按max_submit_job_attempts配置进行指数退避重试(重试延迟计算见 exponential_backoff_retry.py)。此外,该 Executor 支持多团队部署(supports_multi_team = True),并在 Airflow 3.3+ 支持回调执行(supports_callbacks = True),还实现了try_adopt_task_instances以在 Executor 崩溃后按external_executor_id(即 Batch job-id)收养未完成的任务实例。
配置选项(Config Options)
所有配置项既可以写在airflow.cfg的[aws_batch_executor]小节下,也可以通过环境变量以AIRFLOW__AWS_BATCH_EXECUTOR__<OPTION_NAME>的形式设置,例如AIRFLOW__AWS_BATCH_EXECUTOR__JOB_QUEUE = "myJobQueue"。配置节名称与各选项的类型、默认值、示例的权威定义位于 provider 的配置模板中,仓库内对应实现见 get_provider_info.py 与 utils.py 中的CONFIG_GROUP_NAME/CONFIG_DEFAULTS。
必填配置项
| 配置项 | 说明 |
|---|---|
JOB_QUEUE | 作业提交到的 Job Queue,可指定名称或 ARN。必填。 |
JOB_DEFINITION | 作业使用的 Job Definition,可指定名称或 ARN(带或不带修订号,不指定修订号时使用最新激活修订)。必填。 |
JOB_NAME | AWS Batch 作业的名称,最长 128 个字符,首字符必须为字母数字,可含字母、数字、连字符(-)和下划线(_)。必填。 |
REGION_NAME | 部署 Amazon Batch 的 AWS 区域名称,如us-east-1。必填。 |
可选配置项
| 配置项 | 说明 | 默认值 |
|---|---|---|
AWS_CONN_ID | Batch Executor 调用 AWS Batch API 所使用的 Airflow 连接(即凭据)。 | aws_default |
SUBMIT_JOB_KWARGS | JSON 字符串,包含传给 Batchsubmit_jobAPI 的额外参数,例如{"Tags": [{"Key": "key", "Value": "value"}]}。 | 空 |
MAX_SUBMIT_JOB_ATTEMPTS | Batch Executor 尝试提交一个作业的最大次数,针对作业启动失败(API 故障、容器故障等)场景。 | 3 |
CHECK_HEALTH_ON_STARTUP | 是否在启动时检查 Batch Executor 的健康状态。 | True |
配置优先级
当同一配置项在多个位置出现冲突时,从低到高的优先级顺序为:
- 各选项的默认值;
- 通过
airflow.cfg或环境变量显式提供的值(遵循 Airflow 自身的配置优先级机制); SUBMIT_JOB_KWARGS选项中提供的值。
注意:所有运行 Airflow 组件(Scheduler、Webserver、Executor 管理的资源等)的主机/容器上的配置必须保持一致,否则会出现行为不一致。
通过 executor_config 做任务级定制
executor_config是传给 Operator 的可选参数,类型为字典。在 Batch Executor 语境下,它代表一份submit_job_kwargs配置,会被递归地更新到(若存在)Airflow 配置中SUBMIT_JOB_KWARGS之上——近似等价于submit_job_kwargs.update(executor_config),对嵌套字典同样执行递归更新。这使得单个任务可以独立指定 CPU、内存、GPU、环境变量等参数。
源码层面,_submit_job_kwargs()(batch_executor.py)会先深拷贝配置构建的submit_job_kwargs,再用merge_dicts合并executor_config,随后强制写入任务的containerOverrides.command,并自动追加一个环境变量AIRFLOW_IS_EXECUTOR_CONTAINER=true以标记该容器由 Executor 拉起。需要特别留意:executor_config中不允许出现command键,否则execute_async会抛出ValueError。
为 Batch Executor 构建容器镜像
仓库在 executors/Dockerfile 中提供了可直接使用的示例 Dockerfile,构建出的镜像可用于让 AWS Batch 以 Batch Executor 方式运行 Airflow 任务。镜像内置 AWS CLI/API 集成,并支持从 S3 桶或本地文件夹两种方式加载 DAG。
前置条件
构建镜像前需在本机安装 Docker。若容器需要与 AWS 服务交互(例如读取 S3 上的 DAG、写入远程日志),镜像内已安装 AWS CLI,可通过多种方式向容器传递 AWS 认证信息。
构建镜像:方法一(推荐,Iam Role / 默认区域)
使用构建参数aws_default_region构建:
docker build -t my-airflow-image \ --build-arg aws_default_region=YOUR_DEFAULT_REGION .注意镜像的构建与运行需保持同一架构。例如 Apple Silicon 用户可借助docker buildx指定平台:
docker buildx build --platform=linux/amd64 -t my-airflow-image \ --build-arg aws_default_region=YOUR_DEFAULT_REGION .构建镜像:方法二(构建期传入显式凭据)
也可通过构建期参数aws_access_key_id、aws_secret_access_key、aws_default_region、aws_session_token传入 AWS 认证信息:
docker build -t my-airflow-image \ --build-arg aws_access_key_id=YOUR_ACCESS_KEY \ --build-arg aws_secret_access_key=YOUR_SECRET_KEY \ --build-arg aws_default_region=YOUR_DEFAULT_REGION \ --build-arg aws_session_token=YOUR_SESSION_TOKEN .警告:该方法不建议用于生产环境,因为用户凭据会以环境变量的形式固化进镜像,存在安全风险。生产环境应优先使用下文介绍的 IAM 角色方案(由容器运行平台注入临时凭据)。
基础镜像与版本对齐
镜像基于apache/airflow:latest构建。关键要求是:镜像中的 Airflow 与 Python 版本必须与运行 Scheduler 进程(即运行 Executor 的主机/容器)上的 Airflow 与 Python 版本一致。可分别用以下命令校验:
docker run <image_name> version docker run <image_name> python --version例如apache/airflow镜像标签latest-python3.10表示内置 Python 3.10,可按需选用与 Scheduler 环境匹配的带特定 Python 版本的镜像。
加载 DAG
镜像预置了两种 DAG 加载方式(也支持其他自定义方式):
从 S3 桶加载:取消 Dockerfile 中相应 ENTRYPOINT 行的注释,使容器启动时执行aws s3 sync将指定 S3 桶内容同步到容器内/opt/airflow/dags;若想存放到其他目录,可通过container_dag_path构建参数指定。构建时添加--build-arg s3_uri=YOUR_S3_URI,并确保有读取该桶的权限:
docker build -t my-airflow-image \ --build-arg aws_access_key_id=YOUR_ACCESS_KEY \ --build-arg aws_secret_access_key=YOUR_SECRET_KEY \ --build-arg aws_default_region=YOUR_DEFAULT_REGION \ --build-arg aws_session_token=YOUR_SESSION_TOKEN \ --build-arg s3_uri=YOUR_S3_URI .从本地文件夹加载:将 DAG 文件放入 docker build 上下文内的文件夹,通过host_dag_path构建参数指定其位置。默认复制到/opt/airflow/dags,可用container_dag_path参数修改:
docker build -t my-airflow-image --build-arg host_dag_path=./dags_on_host --build-arg container_dag_path=/path/on/container .若将 DAG 加载到/opt/airflow/dags之外的其他路径,需要同步更新 Airflow 配置中对应的 DAG 目录设置(dags_folder)。
安装 Python 依赖
Dockerfile 支持通过pip从requirements.txt安装 Python 依赖。将requirements.txt放在 Dockerfile 同目录(若在其他位置,可用requirements_path构建参数指定,注意 Docker 构建上下文),然后取消 Dockerfile 中复制文件并执行pip install的两行注释即可。
安全最佳实践:使用 IAM 角色
最安全的认证方式是使用 IAM 角色。在 AWS Batch 控制台创建 Job Definition 时,可同时指定Job Role与Execution Role两种角色:
- Execution Role:由容器 agent 使用,代表你向 AWS 发起 API 请求。根据 Batch Executor 所用的计算环境,需要为 Execution Role 附加相应的策略;此外该角色至少需要
CloudWatchLogsFullAccess(或CloudWatchLogsFullAccessV2)策略以写入容器日志。 - Job Role:由容器内的应用进程使用,向 AWS 发起 API 请求。该角色的权限需要基于 DAG 中任务的具体需求来授予;若通过 S3 桶加载 DAG,则该角色需要具备读取该 S3 桶的权限。
创建新 Job Role 或 Execution Role 的步骤:
- 登录 AWS 控制台进入 IAM 页面,在左侧 Access Management 下选择 Roles;
- 在 Roles 页面点击右上角 Create role;
- Trusted entity type 选择 AWS Service;
- 选择适用的 use case;
- 在 Permissions 页面按 Job Role 或 Execution Role 的需求选择权限,完成后点击 Next;
- 输入角色名称与可选描述,检查 Trusted Entities 与权限,按需添加标签,点击 Create role。
创建 Batch Job Definition 时,将上述新建的 Job Role 与 Execution Role 分别选到对应字段即可。若需将 AWS 凭据显式传入容器(如开发环境),可在构建镜像时传入(见上文方法二),但生产环境务必以 IAM 角色替代。
日志:配置远程日志
通过 Batch Executor 运行的任务位于所配置的 VPC 内部,其日志无法被 Airflow UI 直接访问;任务完成后若未做持久化,日志将永久丢失。因此使用 Batch Executor 时必须启用远程日志,以便任务日志持久化并可从 Airflow UI 查看。
配置远程日志时需要注意:
- Airflow 远程日志配置需要在所有运行 Airflow 的主机和容器上保持一致:Webserver 需要该配置以从远程位置拉取日志,Batch Executor 拉起的容器需要该配置以将日志上传到远程位置;
- 将远程日志配置注入容器有多种方式,包括但不限于:
- 在 Dockerfile 中直接以环境变量导出(参见上文 Dockerfile 小节);
- 在 Dockerfile 中更新
airflow.cfg,或复制/挂载/下载一份自定义airflow.cfg; - 在 Job Definition 中以环境变量形式添加;
- 容器内必须配置凭据才能与远程日志服务(如 S3、CloudWatch Logs)交互,常见方式包括:
- 在 Dockerfile 中直接导出凭据;
- 配置一个 Airflow Connection,并将其指定为
remote_log_conn_id(通过上述任意方式注入容器)。Airflow 将仅使用该连接凭据与所选的远程日志目的地交互。
注意:配置项必须在所有运行 Airflow 组件的环境(Scheduler、Webserver、Executor 管理的资源等)中保持一致。
快速上手指南:设置 Batch Executor
让 Batch Executor 在 Apache Airflow 中工作共有 3 个步骤:
- 创建 Airflow 与 Batch 执行任务都能连接到的数据库;
- 创建并配置可运行 Airflow 任务的 Batch 资源;
- 配置 Airflow 使用 Batch Executor 与该数据库。
下文以 AWS 上的 PostgreSQL RDS 实例为数据库后端、EC2 编排类型的 AWS Batch 为例逐步说明。
第一步:创建 RDS DB 实例
在 AWS 控制台进入 RDS 服务,点击 Create database:
- 选择 Standard create,数据库引擎选 PostgreSQL;
- 选择合适的模板、可用性与持久化配置。注意:撰写本文时,"Multi-AZ DB Cluster" 选项不支持设置数据库名称,而数据库名称是后续必需项,因此需避开该选项;
- 设置 DB 实例名、用户名和密码;
- 选择实例配置与存储参数;
- Connectivity 部分选择 "Don't connect to an EC2 compute resource";
- 选择或创建 VPC 与子网,允许对数据库的公共访问;选择或创建安全组与可用区;
- 打开 Additional Configuration 选项卡,将数据库名设为
airflow_db; - 按需选择其他设置,点击 Create database 完成创建。
测试连通性:需要先从你的 IP 地址放行对数据库的入站流量:
- 在 RDS 实例的 "Connectivity & security" 选项卡的 Security 下找到该实例关联的 VPC 安全组;
- 添加入站规则,允许来自你的 IP 地址、TCP 端口 5432(PostgreSQL)的流量;
- 修改安全组后,使用
psql验证连接(需本机安装psql):
psql -h <endpoint> -p 5432 -U <username> <db_name>endpoint 位于 "Connectivity and Security" 选项卡,用户名/密码即创建数据库时设置的凭据,db_name应为airflow_db(除非创建时用了别的名称)。连接成功后会提示输入密码。注意:测试前应确保数据库状态为Available。
第二步:设置 AWS Batch
AWS Batch 有多种编排类型,本指南以 EC2 为例。首先需要构建好上文所述的 Docker 镜像,并将其推送到容器可以拉取到的仓库,这里使用 Amazon Elastic Container Registry(ECR)。
创建 ECR 仓库:
- 进入 ECR 服务,点击 Create repository;
- 命名仓库并按需填写其他信息;
- 点击 Create Repository;
- 创建后进入仓库,点击右上角 "View push commands",按提示将 Docker 镜像推送上去(替换镜像名);推送完成后刷新页面确认镜像已上传。
配置 AWS Batch:
- 登录 AWS 管理控制台,进入 AWS Batch 首页;
- 点击左侧 Wizard,向导将引导创建运行 Batch 作业所需的全部资源;
- 选择编排类型 Amazon EC2;
- 点击 Next。
创建 Compute Environment(计算环境):
- 为计算环境命名、添加标签与合适的实例配置。此处可设置最小、最大与期望 vCPU 数量,以及要使用的 EC2 实例类型;
- Instance Role 选择新建或复用具备所需 IAM 权限的实例配置文件。该实例配置文件允许为计算环境创建的 ECS 容器实例代表你调用所需的 AWS API;
- 选择可访问互联网的 VPC,以及具备必要权限的安全组;
- 点击 Next。
创建 Job Queue(作业队列):
- 为作业队列命名并设置优先级,计算环境选择上一步创建的 Compute Environment。
创建 Job Definition(作业定义):
- 为 Job Definition 命名;
- 选择适当的平台配置,确保启用 "Assign public IP";
- 选择 Execution Role,并确保该角色具备完成任务所需的权限;
- 输入上一步推送到 ECR 的镜像 URI,确保所用角色具备拉取该镜像的权限;
- 选择合适的 Job Role,结合所运行任务的权限需求;
- 按需配置环境,可指定容器可用的 vCPU、内存或 GPU 数量。同时向容器添加以下环境变量:
AIRFLOW__DATABASE__SQL_ALCHEMY_CONN,值为 PostgreSQL 连接串,格式如下(使用上一步创建 RDS 时的值):
postgresql+psycopg2://<username>:<password>@<endpoint>/<database_name>注意:
psycopg2在该 Executor 支持的所有 Airflow 版本上均可用。从 Airflow 3.2.0 起,若安装了psycopg(v3),也可改用postgresql+psycopg://。Airflow 3.2.0 之前的版本不保证 SQLAlchemy 2.0(2.11 与 3.0.x 系列将 SQLAlchemy 钉在<2.0),而 SQLAlchemy 1.4 没有postgresql+psycopg方言——此时 Airflow 会因sqlalchemy.exc.NoSuchModuleError而无法启动。
- 按需添加其他 Airflow 通用配置、Batch Executor 配置(见配置选项小节)或远程日志配置。任何配置变更都应在整个 Airflow 环境中同步,以保持配置一致;
- 点击 Next;
- 在 Review and Create 页面复核所有选择,确认无误后点击 Create Resources。
允许容器访问 RDS 数据库:
最后需要为 Batch 管理的容器配置数据库访问权限。网络配置方式很多,一种可行方案是:
- 登录 AWS 控制台进入 VPC Dashboard;
- 在左侧 Security 下点击 Security groups;
- 选择与 RDS 实例关联的安全组,点击 Edit inbound rules;
- 添加一条规则,允许 PostgreSQL 类型流量来自 Batch Compute Environment 关联子网的 CIDR。
第三步:配置 Airflow
要使用 Batch Executor 并利用上述资源,在运行 Airflow 的环境中定义以下环境变量:
AIRFLOW__CORE__EXECUTOR='airflow.providers.amazon.aws.executors.batch.batch_executor.AwsBatchExecutor' AIRFLOW__DATABASE__SQL_ALCHEMY_CONN=<postgres-connection-string> AIRFLOW__AWS_BATCH_EXECUTOR__REGION_NAME=<executor-region> AIRFLOW__AWS_BATCH_EXECUTOR__JOB_QUEUE=<batch-job-queue> AIRFLOW__AWS_BATCH_EXECUTOR__JOB_DEFINITION=<batch-job-definition> AIRFLOW__AWS_BATCH_EXECUTOR__JOB_NAME=<batch-job-name>初始化 Airflow 数据库:Airflow 数据库在使用前需要初始化,并创建用于登录的用户。下面的命令会创建一个 admin 用户(若数据库尚未初始化,该命令也会一并完成初始化),应在启动 Scheduler 和 Webserver 之前,于运行这两个进程的主机上执行:
airflow users create --username admin --password admin --firstname <your first name> --lastname <your last name> --email <your email> --role Admin远程日志等其他任何配置变更,都建议追加到这个初始化脚本中,以保证 Airflow 环境中各组件配置一致。完成以上三步后,Scheduler 即会通过AwsBatchExecutor把每个任务以独立 AWS Batch 作业的形式投递执行,Airflow 与 Batch 之间的状态同步、失败重试与任务收养均由 Executor 在心跳循环中自动完成。
常见故障排查要点
- 启动即失败:若
CHECK_HEALTH_ON_STARTUP为True,Executor 启动时会用无效 job-id("a"*32)调用describe_jobs做健康检查,任何ClientError或异常都会阻止 Scheduler 启动(见 batch_executor.py),此时应重点核对AWS_CONN_ID、REGION_NAME与 IAM 权限。 - 凭据失效:当
ExpiredTokenException、InvalidClientTokenId、UnrecognizedClientException出现时,Executor 会将连接标记为不健康并按指数退避策略重新加载连接(batch_executor.py),Scheduler 不会崩溃,但需及时修复凭据。 - 作业提交失败:
SUBMIT_JOB_KWARGS中若包含nodeOverrides(多节点作业)或eksPropertiesOverride(EKS 作业),配置加载时会直接抛出KeyError,因为当前实现不支持这两类作业(batch_executor_config.py)。 - 容器内无法访问 AWS 服务:优先检查 Job Definition 中 Execution Role 与 Job Role 的策略是否覆盖容器实际调用的 API(如 S3、CloudWatch Logs)。
参考资料
- Executor 主文档:providers/amazon/docs/executors/batch-executor.rst
- 三个 Executor 共用素材(配置优先级、Dockerfile 说明、日志、RDS、ECR 等):providers/amazon/docs/executors/general.rst
- Executor 核心实现:providers/amazon/src/airflow/providers/amazon/aws/executors/batch/batch_executor.py
- 配置构建与校验:providers/amazon/src/airflow/providers/amazon/aws/executors/batch/batch_executor_config.py
- 配置键、默认值与状态映射:providers/amazon/src/airflow/providers/amazon/aws/executors/batch/utils.py
- API 响应 Schema:providers/amazon/src/airflow/providers/amazon/aws/executors/batch/boto_schema.py
- 示例 Dockerfile:providers/amazon/src/airflow/providers/amazon/aws/executors/Dockerfile
- 配置项类型/默认值/示例的权威定义:providers/amazon/src/airflow/providers/amazon/get_provider_info.py
- 重试延迟计算:providers/amazon/src/airflow/providers/amazon/aws/executors/utils/exponential_backoff_retry.py
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考