MWAA中Airflow获取ECS Task状态及自定义Sensor实现咨询
自定义ECS Task Sensor实现指引
核心结论
目前Airflow官方Amazon Provider确实未提供开箱即用的ECS Task状态Sensor,但自定义实现成本很低:不需要从零开发自定义ECS Hook。MWAA预装的Amazon Airflow Provider包已经自带成熟的EcsHook,封装了ECS API鉴权、客户端初始化、区域适配等所有通用逻辑,你不需要参考EMR Hook的实现重写Hook层,直接在自定义Sensor中复用现有EcsHook即可,实现链路比EMR的状态感知更简单。
ECS状态与Airflow状态映射规则
你触发ECS Task后会拿到唯一的task_arn,ECS返回的任务状态可以直接按以下规则映射,不需要额外做二次判断:
- 成功终止态:任务状态为
STOPPED,且任务下所有容器的退出码均为0 → Sensor返回成功,标记Airflow任务完成 - 失败终止态:任务状态为
STOPPED,且存在容器退出码非0、任务被强制终止、资源调度失败(比如配额不足、子网/安全组配置错误、找不到匹配的EC2实例)→ Sensor抛出AirflowException,标记Airflow任务失败 - 中间运行态:任务状态为
PROVISIONING/PENDING/ACTIVATING/RUNNING/DEACTIVATING/STOPPING→ Sensor返回False,按照设定的间隔继续轮询
具体实现代码
自定义Sensor代码
直接继承Airflow内置BaseSensor实现即可,核心逻辑就是调用EcsHook封装的boto3客户端轮询任务状态:
from airflow.sensors.base import BaseSensorOperator from airflow.providers.amazon.aws.hooks.ecs import EcsHook from airflow.exceptions import AirflowException class EcsTaskStateSensor(BaseSensorOperator): # 支持模板化传参,方便从XCOM拉取上游触发任务返回的task_arn template_fields = ("task_arn", "cluster") def __init__( self, task_arn: str, cluster: str, aws_conn_id: str = "aws_default", **kwargs ): super().__init__(**kwargs) self.task_arn = task_arn self.cluster = cluster self.aws_conn_id = aws_conn_id def poke(self, context): hook = EcsHook(aws_conn_id=self.aws_conn_id) resp = hook.conn.describe_tasks( cluster=self.cluster, tasks=[self.task_arn] ) if not resp.get("tasks"): raise AirflowException(f"ECS task {self.task_arn} not found in cluster {self.cluster}") task_detail = resp["tasks"][0] last_status = task_detail["lastStatus"] self.log.info(f"ECS task {self.task_arn} current status: {last_status}") if last_status != "STOPPED": return False # 校验STOPPED状态下的容器退出码 containers = task_detail.get("containers", []) for container in containers: exit_code = container.get("exitCode") if exit_code != 0: stop_reason = task_detail.get("stoppedReason", "No stop reason provided") raise AirflowException( f"ECS task {self.task_arn} failed, stop reason: {stop_reason}, " f"container {container['name']} exit code: {exit_code}" ) self.log.info(f"ECS task {self.task_arn} completed successfully") return True
DAG中调用示例
和EMR Operator+Sensor的使用逻辑完全一致,上游触发ECS任务后通过XCOM传递task_arn给下游Sensor即可:
from airflow import DAG from airflow.providers.amazon.aws.operators.ecs import EcsRunTaskOperator from datetime import datetime with DAG( dag_id="ecs_task_workflow", start_date=datetime(2024, 1, 1), schedule=None, catchup=False ) as dag: trigger_ecs_task = EcsRunTaskOperator( task_id="trigger_ecs_task", task_definition="your-pre-created-task-def-name", cluster="your-ecs-cluster-name", launch_type="FARGATE", # 按实际环境补充子网、安全组、执行角色等配置 awslogs_group="/ecs/your-log-group", awslogs_stream_prefix="ecs", do_xcom_push=True ) wait_for_task_done = EcsTaskStateSensor( task_id="wait_for_task_done", task_arn="{{ ti.xcom_pull(task_ids='trigger_ecs_task', key='return_value') }}", cluster="your-ecs-cluster-name", mode="reschedule", poke_interval=30, timeout=3600 ) trigger_ecs_task >> wait_for_task_done
落地注意事项
- 先确认MWAA环境对应的Amazon Provider版本:版本号≥2.0即可直接使用内置
EcsHook,如果是低于2.0的老版本MWAA,把对应版本的Amazon Provider whl包上传到DAG目录下的plugins路径即可,不需要修改环境配置。 - 上述状态判断逻辑对Fargate、EC2启动类型的ECS Task完全通用,不需要针对启动类型做额外适配。
- 轮询模式优先选
reschedule:如果你的ECS任务平均运行时长超过10分钟,reschedule模式会在轮询间隔释放Worker资源,避免长期占用MWAA的Worker并发槽位,比默认的poke模式资源利用率更高。 poke_interval和timeout参数按需调整:常规短任务设30秒轮询、1小时超时即可,长时任务可以把轮询间隔调到1-2分钟,超时时间匹配任务的SLA上限即可。
内容的提问来源于stack exchange,提问作者Gopi Gof
相关产品推荐
相关产品推荐

