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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 00:01:05