Apache Airflow AWS ECS Executor 实战指南:架构原理、配置调优与 Fargate 部署全流程

Apache Airflow AWS ECS Executor 实战指南:架构原理、配置调优与 Fargate 部署全流程 Apache Airflow AWS ECS Executor 实战指南架构原理、配置调优与 Fargate 部署全流程【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache Airflow 的 AWS ECS Executor 将调度器生成的每一个任务以独立 ECS 容器的方式运行在 Amazon Elastic Container Service 之上让每个 Airflow 任务都拥有完全隔离的 CPU、内存与网络资源同时按需启动、按运行时长计费无需长期驻留的 worker 节点。本文基于 Apache Airflow 仓库中 ECS Executor 官方文档 展开并结合 ecs_executor.py 等源码实现系统讲解其工作原理、全部配置项、任务镜像构建、IAM 权限、日志体系与性能调优最后给出从 RDS 数据库、ECS Fargate 集群到 Airflow 配置的端到端部署步骤帮助读者在生产环境落地这一按需弹性执行方案。为什么选择 ECS Executor三大核心优势ECS Executor 是 Apache Airflow 众多执行器中的一种其核心思想是Airflow 调度器为每个任务生成执行命令后通过 ECS API 将命令投递到远程 ECS 集群中的独立容器内运行调度器再周期性地查询这些 ECS 任务的运行状态。官方文档总结了这种模式带来的三个关键收益任务隔离Task Isolation没有共享的 worker任何一个任务的 CPU、内存、磁盘资源都被独立隔离不会出现吵闹邻居效应某个任务引发的网络故障或整个容器崩溃只会影响它自身单个用户也无法通过并发触发大量任务而压垮共享环境。自定义环境Customized Environments可以为不同任务构建包含特定系统级依赖、二进制文件或数据的不同容器镜像环境能力随镜像而定不再受限于统一的 worker 环境。成本效益Cost Effective计算资源只存在于 Airflow 任务运行的生命周期内无需为随时待命的长期 worker 付费也省去了对这些常驻节点的维护与打补丁成本。从源码来看AwsEcsExecutor继承自BaseExecutor其类注释明确描述了这一工作模式调度器创建 shell 命令并交给执行器执行器在远程 ECS 集群上用与调度器相同镜像的 task definition 启动容器随后通过任务 ARN 周期性轮询状态见 ecs_executor.py。理解执行流程从调度到容器运行的源码级原理要正确使用 ECS Executor先要理解它的运行闭环。结合 ecs_executor.py 的源码核心流程可以拆解为以下几个环节入队execute_async调度器将任务交给执行器后execute_async()会把任务包装成EcsQueuedTask放入pending_workloads队列。值得注意的是如果executor_config中出现了name或command键会直接抛出ValueError因为这两个字段由执行器接管不允许用户覆盖ecs_executor.py。心跳同步sync执行器的sync()每轮心跳先确保 boto3 连接健康不健康时按指数退避重连再依次调用sync_running_workloads()轮询在跑任务、attempt_workload_runs()尝试启动新任务。若检测到 AWS 凭据失效ExpiredTokenException、InvalidClientTokenId、UnrecognizedClientException等会将连接标记为不健康并记录告警后重试而不是让异常把调度器进程打死ecs_executor.py。启动任务_run_task_run_task()将 Airflow 命令写入 ECS API 的containerOverrides[0].command调用 boto3 的ecs.run_task(**kwargs)启动容器并使用BotoRunTaskSchema校验响应ecs_executor.py。状态轮询describe_tasks由于 AWS 的describe_tasks接口单次最多接收 100 个 ARN源码将批量查询大小固定为DESCRIBE_TASKS_BATCH_SIZE 99分批查询所有在跑任务ecs_executor.py。状态判定get_task_stateECS 任务被映射为 Airflow 的QUEUED / RUNNING / REMOVED / FAILED / SUCCESS五种状态——RUNNING状态对应运行中desired_status为RUNNING而尚未启动视为排队任务超时未启动会判定为REMOVED容器全部正常退出exit code 全为 0判定为SUCCESS否则为FAILED见 utils.py。失败重试无论是 API 调用失败还是容器启动失败任务都会被重新放回pending_workloads在未超过max_run_task_attempts前按指数退避延迟重试超过上限后标记为 FAILED 并记录ecs task submit failure事件ecs_executor.py。优雅退出与接管调度器收到 SIGTERM 时terminate()会调用stop_task停止所有在跑 ECS 任务新调度器启动后则通过try_adopt_task_instances()依据external_executor_id即 ECS 任务 ARN接管仍在运行的容器避免任务中断ecs_executor.py。此外执行器启动时还会注入环境变量AIRFLOW_IS_EXECUTOR_CONTAINERtrue用于标记当前进程运行在由执行器拉起的容器中供容器内 Airflow 组件识别运行环境ecs_executor.py。配置选项详解ECS Executor 的全部配置项均位于airflow.cfg的[aws_ecs_executor]小节也可以通过环境变量AIRFLOW__AWS_ECS_EXECUTOR__OPTION_NAME设置例如export AIRFLOW__AWS_ECS_EXECUTOR__CONTAINER_NAMEmyEcsContainer必需配置项配置项说明CLUSTERAmazon ECS 集群的名称必填。CONTAINER_NAME用于执行 Airflow 任务的容器名称该容器必须在 ECS Task Definition 中声明必填。REGION_NAMEECS 所在的 AWS 区域名称必填。可选配置项配置项默认值说明ASSIGN_PUBLIC_IPFalse是否为执行器启动的容器分配公网 IP。AWS_CONN_IDaws_default执行器调用 AWS ECS API 时使用的 Airflow 连接凭据。LAUNCH_TYPEFARGATE启动类型取值为FARGATE或EC2。PLATFORM_VERSIONLATEST使用 Fargate 启动类型时的平台版本。RUN_TASK_KWARGS空传给 ECSrun_taskAPI 的 JSON 字符串参数。SECURITY_GROUPSVPC 默认与 ECS 任务关联的安全组 ID最多 5 个逗号分隔。SUBNETSVPC 默认与 ECS 任务关联的子网 ID最多 16 个逗号分隔。TASK_DEFINITION最新 ACTIVE 修订要运行的 Task Definition格式为family:revision或完整 ARN。MAX_RUN_TASK_ATTEMPTS3任务启动失败如 ECS API 失败、容器失败时的最大重试次数。CHECK_HEALTH_ON_STARTUPTrue是否在启动时对 ECS Executor 执行健康检查。其中MAX_RUN_TASK_ATTEMPTS、CHECK_HEALTH_ON_STARTUP、AWS_CONN_ID、ASSIGN_PUBLIC_IP、PLATFORM_VERSION的默认值可以直接在源码的CONFIG_DEFAULTS字典中得到印证见 utils.py。配置优先级当多个来源同时提供配置时从低到高的生效顺序为各选项的内置默认值通过airflow.cfg或环境变量显式提供的值按 Airflow 标准配置优先级解析RUN_TASK_KWARGS选项中提供的值优先级最高。任务级覆盖executor_config 与 exec_configexecutor_config部分文档与旧版本中写作exec_config是传给算子的可选字典参数在 ECS Executor 语境下它表示一份run_task_kwargs配置会以递归合并的方式覆盖上述第 3 层配置。合并逻辑可以近似理解为run_task_kwargs.update(executor_config) # 递归地对每个嵌套字典执行 update在源码中这一逻辑由merge_dicts(run_task_kwargs, exec_config)实现ecs_executor.py。借助它你可以在 DAG 级别为单个任务指定 CPU、内存、GPU、环境变量等ContainerOverride参数实现任务粒度的资源定制。底层参数组装逻辑配置在进入 boto3 前会经过 ecs_executor_config.py 中的build_task_kwargs()统一加工几个值得注意的实现细节capacity_provider_strategy与launch_type互斥二者同时提供会直接抛出ValueError如果两者都未提供但集群存在默认 capacity provider则使用它否则回退到FARGATE启动类型ecs_executor_config.py。EC2 启动类型自动清理参数使用EC2启动类型时不允许提供platform_version执行器会替用户静默丢弃该参数而非报错ecs_executor_config.py。强制约束count恒为 1容器覆盖列表恒以CONTAINER_NAME指定的容器为第一个且其command由执行器在执行期覆写ecs_executor_config.py。网络配置组装SUBNETS、SECURITY_GROUPS会按逗号拆分填入awsvpcConfigurationASSIGN_PUBLIC_IP被转换为ENABLED/DISABLEDEC2 类型下该项被忽略若提供了子网/安全组但缺少子网会报错At least one subnet is required to run a taskecs_executor_config.py。键名统一 camelCase组装完成后所有键递归转换为 boto3 期望的 camelCase 形式并校验整个字典可被 JSON 序列化ecs_executor_config.py。注意所有配置选项必须在运行 Airflow 各组件的所有主机/环境中保持一致调度器、Webserver、ECS 任务容器等否则可能出现行为不一致。构建 ECS Executor 任务容器镜像仓库在 executors/Dockerfile 提供了一个可直接使用的示例 Dockerfile。该镜像基于apache/airflow:latest构建内置 AWS CLI支持与 AWS 服务交互并预置了从 S3 桶或本地目录加载 DAG 两种方式。构建前置条件构建镜像需要本机安装 Docker。镜像版本对齐重要镜像中的 Airflow 与 Python 版本必须与运行调度器即运行执行器的主机/容器保持一致。构建后可在本地用以下命令核对docker run image_name version docker run image_name python --versionApache Airflow 官方镜像按 Python 版本提供不同 tag例如latest-python3.10表示镜像内置 Python 3.10。方式一通过默认区域构建推荐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四个构建参数也可以直接传入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 .警告该方式会将用户凭据直接写入容器镜像存在安全风险不建议在生产环境使用。加载 DAG 的两种预置方式从 S3 桶加载取消 Dockerfile 中 entrypoint 相关行的注释使构建时把指定 S3 桶的 DAG 同步到容器内的/opt/airflow/dags可通过container_dag_path构建参数改到其他目录构建时传入--build-arg s3_uriYOUR_S3_URI并确保具备桶的读取权限。从本地目录加载将 DAG 文件放在 Docker 构建上下文内的文件夹中通过host_dag_path指定其位置默认拷贝到/opt/airflow/dags可用container_dag_path修改目标路径但若改了路径需要同步更新 Airflow 配置中的 DAG 目录docker build -t my-airflow-image --build-arg host_dag_path./dags_on_host --build-arg container_dag_path/path/on/container .安装 Python 依赖镜像支持通过requirements.txt安装 Python 依赖把requirements.txt放在 Dockerfile 同目录其他位置可用requirements_path构建参数指定并取消 Dockerfile 中拷贝该文件与执行pip install两行代码的注释即可。注意拷贝时要符合 Docker 构建上下文的范围。IAM 角色与权限配置最安全的方式是使用 IAM 角色。创建 ECS Task Definition 时可以指定两类角色Task Execution Role任务执行角色供容器代理代你发起 AWS API 请求。对 ECS Executor 而言该角色至少需要AmazonECSTaskExecutionRolePolicy与CloudWatchLogsFullAccess或CloudWatchLogsFullAccessV2策略。Task Role任务角色供容器内运行的业务代码发起 AWS API 请求其权限取决于 DAG 中定义的任务内容若通过 S3 桶加载 DAG该角色还需具备读取该 S3 桶的权限。在 AWS 控制台创建角色的步骤进入 IAM 页面在左侧 Access Management 下选择 Roles点击右上角 Create roleTrusted entity type 选择 AWS ServiceUse case 下拉框选择 Elastic Container Service具体用例选择Elastic Container Service Task点击 Next在 Permissions 页面按角色类型Task Role 或 Task Execution Role勾选所需权限点击 Next为角色命名并填写可选描述核对 Trusted Entities 与权限按需添加标签点击 Create role。创建 Task Definition 时将这两个新建角色分别选为任务的 Task Role 与 Task Execution Role。日志体系远程日志与 ECS Task 日志为什么必须配置远程日志通过 ECS Executor 运行的任务在配置的 VPC 内执行任务完成后日志对 Airflow UI 不可见且永久丢失。因此官方强烈建议配合远程日志使用将任务日志持久化并可在 Airflow UI 中查看。S3 远程日志与 CloudWatch 远程日志分别由 Amazon 提供商的对应 task handler 支持。远程日志配置需要满足两个要求配置一致性远程日志配置必须同时作用于所有运行 Airflow 的主机与容器——Webserver 需要该配置以便从远端拉取日志展示ECS 任务容器需要该配置以便把日志上传到远端。容器内凭据容器内必须配置可与远端日志服务S3、CloudWatch Logs 等交互的凭据。注入配置与凭据的多种途径把远程日志配置注入容器的方式包括不限于直接在 Dockerfile 中以环境变量导出见上文 Dockerfile 部分在 Dockerfile 中更新airflow.cfg或拷贝/挂载/下载自定义airflow.cfg在 ECS Task Definition 中以明文或通过 AWS Secrets Manager / Systems Manager 的 Secrets 机制注入环境变量使用 ECS Task Environment Files任务环境文件注入。凭据注入方式包括直接在 Dockerfile 中导出凭据配置一个 Airflow Connection并将其指定为远程日志的remote log conn id——Airflow 会专门使用该连接中的凭据与远程日志目标交互。ECS Task 自身日志调试利器ECS 本身可以配置awslogs日志驱动把 ECS 任务进程的全部输出包括airflow tasks run ...的日志与容器生命周期内的其他日志发送到 CloudWatch Logs这在排查远程日志问题或验证远程日志配置时非常有用。注意这部分 ECS 任务日志不会出现在 Airflow Webserver UI 中只能从 CloudWatch Logs 查看。性能与扩展性调优官方文档披露了实测数据由于容器启动耗时ECS Executor 会给每个 Airflow 任务增加约5060 秒的延迟但换来的是更高的并行度与隔离性。团队曾以超过 1,000 个任务并行调度进行测试观察到最多约500 个任务同时运行——该上限与 ECS Service Quotas 一致需要时可向 AWS 申请提升配额。大规模运行 ECS Executor 及 Airflow 时建议关注以下配置项它们会限制并发任务数或影响调度器性能core.max_active_tasks_per_dag单个 DAG 的最大活跃任务数core.max_active_runs_per_dag单个 DAG 的最大活跃运行数core.parallelismAirflow 全局最大并行任务数scheduler.max_tis_per_query调度器单次查询最多获取的任务实例数default_pool_task_slot_count默认任务池的槽位数量scheduler_health_check_threshold调度器健康检查阈值。端到端部署指南三步搭建 ECS Executor官方快速开始指南将整套部署归纳为三步创建数据库 → 创建并配置 ECS 集群 → 配置 Airflow 使用 ECS Executor。下面以 PostgreSQL RDS 与 Fargate 为例给出完整步骤。第一步创建 RDS PostgreSQL 数据库Airflow 与 ECS 中运行的任务需要连接同一个元数据库Airflow 支持多种数据库后端本指南使用 AWS RDS PostgreSQL。登录 AWS 管理控制台进入 RDS 服务点击 Create database 开始创建选择 Standard create数据库引擎选 PostgreSQL按需选择模板、可用性与持久化配置注意Multi-AZ DBCluster选项目前不支持设置数据库名而数据库名是后续必需步骤故不建议选用设置数据库实例名、用户名和密码选择实例规格与存储参数Connectivity 部分选择Dont connect to an EC2 compute resource选择或创建 VPC 与子网允许对数据库的公网访问选择或创建安全组并选择可用区在 Additional Configuration 中把数据库名设为airflow_db按需配置其余设置点击 Create database。测试连通性先在 RDS 实例的 Connectivity security 页签中找到 VPC 安全组链接为其添加入站规则放行你的 IP 在 TCP 5432 端口PostgreSQL的流量待数据库状态为Available后用psql验证psql -h endpoint -p 5432 -U username db_nameendpoint 可在 Connectivity and Security 页签找到用户名/密码即创建时的凭据db_name为airflow_db除非创建时另设。连接成功后系统会提示输入密码。第二步创建 ECS Fargate 集群与 Task Definition推送镜像到 ECR镜像构建完成后需要放入 ECS 可拉取的仓库本指南使用 Amazon ECR进入 ECR 服务点击 Create repository为仓库命名并按需填写信息点击创建进入仓库点击右上角 View push commands按提示推送 Docker 镜像替换镜像名推送后刷新页面确认镜像已上传。创建 ECS 集群登录 AWS 管理控制台进入 Amazon Elastic Container Service点击 Clusters → Create ClusterInfrastructure 下确保选中AWS Fargate (Serverless)按需选择其他选项点击 Create 创建集群。创建 Task Definition左侧点击 Task Definitions → Create new task definition设置 Task Definition Family 名称Launch Type 选择AWS Fargate选择或创建 Task Role 与 Task Execution Role确保角色具备完成各自任务所需的权限可直接创建只含最小必需权限的 Task Execution Role为容器命名镜像 URI 填入上一步推送到 ECR 的镜像地址并确保所用角色具备拉取该镜像的权限为容器添加如下环境变量AIRFLOW__DATABASE__SQL_ALCHEMY_CONN值使用上一步数据库配置格式为postgresqlpsycopg://username:passwordendpoint/database_nameAIRFLOW__ECS_EXECUTOR__SECURITY_GROUPS值为与 RDS 实例所在 VPC 关联的、逗号分隔的安全组 ID 列表AIRFLOW__ECS_EXECUTOR__SUBNETS值为与 RDS 实例关联的、逗号分隔的子网 ID 列表按需添加 Airflow 通用配置、ECS Executor 配置见配置选项详解或远程日志配置见日志体系。注意任何配置变更都应在整个 Airflow 环境中同步以保持配置一致点击 Create 完成创建。允许 ECS 容器访问 RDS 数据库可选的一种网络方案进入 VPC Dashboard左侧 Security 下选择 Security groups选中与 RDS 实例关联的安全组点击 Edit inbound rules新增一条规则允许 PostgreSQL 类型流量访问 ECS 集群所在子网段的 CIDR。第三步配置 Airflow 使用 ECS Executor创建一个脚本例如ecs_executor_config.sh内容如下export AIRFLOW__CORE__EXECUTORairflow.providers.amazon.aws.executors.ecs.ecs_executor.AwsEcsExecutor export AIRFLOW__DATABASE__SQL_ALCHEMY_CONNpostgres-connection-string export AIRFLOW__AWS_ECS_EXECUTOR__REGION_NAMEexecutor-region export AIRFLOW__AWS_ECS_EXECUTOR__CLUSTERecs-cluster-name export AIRFLOW__AWS_ECS_EXECUTOR__CONTAINER_NAMEecs-container-name export AIRFLOW__AWS_ECS_EXECUTOR__TASK_DEFINITIONtask-definition-name export AIRFLOW__AWS_ECS_EXECUTOR__LAUNCH_TYPEFARGATE export AIRFLOW__AWS_ECS_EXECUTOR__PLATFORM_VERSIONLATEST export AIRFLOW__AWS_ECS_EXECUTOR__ASSIGN_PUBLIC_IPTrue export AIRFLOW__AWS_ECS_EXECUTOR__SECURITY_GROUPSsecurity-group-id-for-rds export AIRFLOW__AWS_ECS_EXECUTOR__SUBNETSsubnet-id-for-rds其中AIRFLOW__CORE__EXECUTOR指向的AwsEcsExecutor类正是 ecs_executor.py 中定义的执行器实现类。该脚本需要在启动 Airflow Scheduler 与 Webserver 之前于对应主机上执行远程日志等其他配置变更也应追加到该脚本中保证整个 Airflow 环境配置一致。初始化 Airflow 数据库并创建用户Airflow 元数据库在首次使用前需要初始化并添加一个可登录的用户。下面命令创建管理员用户若数据库尚未初始化该命令会一并完成初始化airflow users create --username admin --password admin --firstname your first name --lastname your last name --email your email --role Admin验证与进一步探索完成上述步骤后Scheduler 会以 ECS Executor 身份运行每个被调度的任务都会在 ECS 集群中以独立容器方式启动。执行器启动时的健康检查会在CHECK_HEALTH_ON_STARTUPTrue默认开启时执行它故意用一个 32 位无效任务 ID 调用stop_task只有当 ECS 返回task was not found的InvalidParameterException时才判定连接健康否则直接阻止 Scheduler 启动ecs_executor.py这一机制可帮助你在第一时间发现权限、集群或区域配置问题。仓库中配套的单元测试如 test_ecs_executor.py、test_utils.py覆盖了入队校验、任务状态迁移、失败重试、配置组装等关键路径可作为理解执行器行为细节的补充阅读材料。若需为任务定制资源可参考 DAG 中为算子传入executor_config对应 ECSContainerOverride结构的方式实现 CPU、内存、GPU 与额外环境变量的任务级控制。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考