资讯动态

Apache Airflow AWS Batch Executor 实战指南:以 Amazon Batch 弹性运行工作流

发布时间:2026/9/13 23:59:39 来源:尧图企业网站定制
Apache Airflow AWS Batch Executor 实战指南以 Amazon Batch 弹性运行工作流【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache Airflow 的 AWS Batch ExecutorAwsBatchExecutor将 Scheduler 派发的每一个任务封装为一次 AWS Batch 作业在由 Amazon Batch 动态供给的独立容器中执行从而把任务调度与资源供给彻底解耦。本指南围绕该 Executor 的配置项、容器镜像构建、IAM 权限、远程日志与端到端快速上手指南展开结合 amazon provider 源码说明其底层调用链读完即可在自己的 Airflow 环境中落地一套可弹性伸缩、按需付费的批处理执行体系。AWS Batch Executor 是什么AWS Batch Executor 是 Apache Airflow 的一种 Executor 实现位于 amazon provider 包中核心类为AwsBatchExecutorbatch_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 支持 FargateECS、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.pySUBMITTED/PENDING/RUNNABLE/STARTING映射为 Airflow 的QUEUEDRUNNING映射为RUNNINGSUCCEEDED映射为SUCCESSFAILED映射为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_NAMEAWS Batch 作业的名称最长 128 个字符首字符必须为字母数字可含字母、数字、连字符-和下划线_。必填。REGION_NAME部署 Amazon Batch 的 AWS 区域名称如us-east-1。必填。可选配置项配置项说明默认值AWS_CONN_IDBatch Executor 调用 AWS Batch API 所使用的 Airflow 连接即凭据。aws_defaultSUBMIT_JOB_KWARGSJSON 字符串包含传给 Batchsubmit_jobAPI 的额外参数例如{Tags: [{Key: key, Value: value}]}。空MAX_SUBMIT_JOB_ATTEMPTSBatch Executor 尝试提交一个作业的最大次数针对作业启动失败API 故障、容器故障等场景。3CHECK_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_CONTAINERtrue以标记该容器由 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_regionYOUR_DEFAULT_REGION .注意镜像的构建与运行需保持同一架构。例如 Apple Silicon 用户可借助docker buildx指定平台docker buildx build --platformlinux/amd64 -t my-airflow-image \ --build-arg aws_default_regionYOUR_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_idYOUR_ACCESS_KEY \ --build-arg aws_secret_access_keyYOUR_SECRET_KEY \ --build-arg aws_default_regionYOUR_DEFAULT_REGION \ --build-arg aws_session_tokenYOUR_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_uriYOUR_S3_URI并确保有读取该桶的权限docker build -t my-airflow-image \ --build-arg aws_access_key_idYOUR_ACCESS_KEY \ --build-arg aws_secret_access_keyYOUR_SECRET_KEY \ --build-arg aws_default_regionYOUR_DEFAULT_REGION \ --build-arg aws_session_tokenYOUR_SESSION_TOKEN \ --build-arg s3_uriYOUR_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 roleTrusted 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 部分选择 Dont connect to an EC2 compute resource选择或创建 VPC 与子网允许对数据库的公共访问选择或创建安全组与可用区打开 Additional Configuration 选项卡将数据库名设为airflow_db按需选择其他设置点击 Create database 完成创建。测试连通性需要先从你的 IP 地址放行对数据库的入站流量在 RDS 实例的 Connectivity security 选项卡的 Security 下找到该实例关联的 VPC 安全组添加入站规则允许来自你的 IP 地址、TCP 端口 5432PostgreSQL的流量修改安全组后使用psql验证连接需本机安装psqlpsql -h endpoint -p 5432 -U username db_nameendpoint 位于 Connectivity and Security 选项卡用户名/密码即创建数据库时设置的凭据db_name应为airflow_db除非创建时用了别的名称。连接成功后会提示输入密码。注意测试前应确保数据库状态为Available。第二步设置 AWS BatchAWS Batch 有多种编排类型本指南以 EC2 为例。首先需要构建好上文所述的 Docker 镜像并将其推送到容器可以拉取到的仓库这里使用 Amazon Elastic Container RegistryECR。创建 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 时的值postgresqlpsycopg2://username:passwordendpoint/database_name注意psycopg2在该 Executor 支持的所有 Airflow 版本上均可用。从 Airflow 3.2.0 起若安装了psycopgv3也可改用postgresqlpsycopg://。Airflow 3.2.0 之前的版本不保证 SQLAlchemy 2.02.11 与 3.0.x 系列将 SQLAlchemy 钉在2.0而 SQLAlchemy 1.4 没有postgresqlpsycopg方言——此时 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__EXECUTORairflow.providers.amazon.aws.executors.batch.batch_executor.AwsBatchExecutor AIRFLOW__DATABASE__SQL_ALCHEMY_CONNpostgres-connection-string AIRFLOW__AWS_BATCH_EXECUTOR__REGION_NAMEexecutor-region AIRFLOW__AWS_BATCH_EXECUTOR__JOB_QUEUEbatch-job-queue AIRFLOW__AWS_BATCH_EXECUTOR__JOB_DEFINITIONbatch-job-definition AIRFLOW__AWS_BATCH_EXECUTOR__JOB_NAMEbatch-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为TrueExecutor 启动时会用无效 job-ida*32调用describe_jobs做健康检查任何ClientError或异常都会阻止 Scheduler 启动见 batch_executor.py此时应重点核对AWS_CONN_ID、REGION_NAME与 IAM 权限。凭据失效当ExpiredTokenException、InvalidClientTokenId、UnrecognizedClientException出现时Executor 会将连接标记为不健康并按指数退避策略重新加载连接batch_executor.pyScheduler 不会崩溃但需及时修复凭据。作业提交失败SUBMIT_JOB_KWARGS中若包含nodeOverrides多节点作业或eksPropertiesOverrideEKS 作业配置加载时会直接抛出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.rstExecutor 核心实现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.pyAPI 响应 Schemaproviders/amazon/src/airflow/providers/amazon/aws/executors/batch/boto_schema.py示例 Dockerfileproviders/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),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价