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

MWAA中ECSOperator触发ECS任务后,PythonSensor获取Task ARN失败求助

问题解决:Airflow ECSOperator XCom为空无法获取Task ARN

问题背景

使用Airflow 2.2.2 + apache-airflow-providers-amazon 2.4.0,通过ECSOperator触发Fargate任务,开启do_xcom_push=True后,PythonSensor无法从XCom获取Task ARN(XCom始终为空),导致轮询失败。ECS任务本身可正常执行,且因版本限制无法使用ECSTaskStateSensor,也不想升级Airflow版本。

核心问题分析

  1. 任务依赖缺失:原DAG中run_ecs_task与await_task_finish未设置依赖,导致Sensor与ECSOperator并行执行,此时ECSOperator尚未推送XCom,故拉取结果为空。
  2. Task ARN提取路径错误:boto3的run_task响应中,Task ARN的键为taskArn(驼峰命名),而非代码中使用的task_arn。
  3. 空值未处理:未对XCom为空或Task ARN不存在的情况做容错处理,直接抛出异常导致Sensor失败。

解决方案

1. 修复任务依赖

确保ECSOperator执行完成后再启动PythonSensor,修改DAG的任务流向:

run_ecs_task >> await_task_finish >> sns_pass

2. 修正Task ARN提取逻辑

从XCom返回的run_task响应中,正确提取Task ARN:

task_arn = xcom['tasks'][0]['taskArn']

3. 增加空值容错判断

在PythonSensor的可调用函数中,增加XCom未就绪或Task ARN不存在的判断,返回False让Sensor继续轮询,而非直接抛出异常。

修正后的完整代码

import json
import boto3
import os
from datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.amazon.aws.operators.ecs import ECSOperator
from airflow.providers.amazon.aws.operators.sns import SnsPublishOperator
from airflow.sensors.python import PythonSensor

def check_task_completed(task_id, **kwargs):
    task_instance = kwargs['ti']
    xcom = task_instance.xcom_pull(task_ids=task_id)
    
    # 检查XCom是否就绪或是否包含Task信息
    if not xcom or 'tasks' not in xcom or len(xcom['tasks']) == 0:
        print("XCom数据未就绪,等待ECS任务启动完成")
        return False
    
    task_arn = xcom['tasks'][0]['taskArn']
    print(f'Task ARN: {task_arn}')

    client = boto3.client('ecs')
    try:
        response = client.describe_tasks(
            cluster='integration-cluster',
            tasks=[task_arn]
        )
    except client.exceptions.ClientError as e:
        print(f"查询任务状态失败: {str(e)}")
        return False

    if not response.get('tasks'):
        print("未查询到任务状态,继续等待")
        return False

    task = response['tasks'][0]
    container = task['containers'][0]
    if task['lastStatus'] == 'STOPPED' and container['lastStatus'] == 'STOPPED':
        if container['exitCode'] == 0:
            return True
        else:
            raise Exception(f'Task {task_arn} 执行失败: {task["stoppedReason"]}')
    else: 
        return False

def fail_sns(messagetext, fail_sns_arn, context):
    sns = boto3.client('sns')
    sns.publish(
        TopicArn=fail_sns_arn,   
        Message=messagetext,
        Subject=messagetext   
    )
    
DEFAULT_ARGS={
    'owner': 'airflow',
    'depends_on_past': False,
    'email': ['airflow@example.com'],
    'email_on_failure': False,
    'email_on_retry': False
}

DAG_ID=os.path.basename(__file__).replace('.py', '')
with DAG(
    tags=["tag1", "tag2", "tag3"],
    dag_id=DAG_ID,
    default_args=DEFAULT_ARGS,
    dagrun_timeout=timedelta(hours=4),
    start_date=datetime(2023, 5, 4, 0 , 0), 
    schedule_interval='0 2 * * *' 
) as dag:

    success_sns_message='CLIENTNAME_SYSTEMNAME_SUCCESS'
    fail_sns_message='CLIENTNAME_SYSTEMNAME_FAIL'
    success_sns_arn='arn:aws:sns:us-east-1:000000000000:CLIENTNAME_SYSTEMNAME_SUCCESS'
    fail_sns_arn='arn:aws:sns:us-east-1:000000000000:CLIENTNAME_SYSTEMNAME_FAIL'

    run_ecs_task=ECSOperator(
        task_id="ecs_operator_task",
        dag=dag,
        cluster="cluster-name",
        task_definition="task-definition-name",
        launch_type="FARGATE",
        overrides={
            "containerOverrides": [
                {
                    "name": "SYSTEMNAME-integration",
                    "command": [
                        "CLIENTNAME"
                    ],
                },
            ],
        },
        network_configuration={
            "awsvpcConfiguration": {
                "subnets": ["subnet-00000000", "subnet-00000001"],
                "assignPublicIp": "ENABLED",
                "securityGroups": ["sg-00000000000000000"],
            }
        },
        awslogs_group="/ecs/SYSTEMNAME-integration-task-definition",
        awslogs_stream_prefix="ecs/SYSTEMNAME-integration",
        do_xcom_push=True
    )

    await_task_finish=PythonSensor(
        task_id="py_sensor_task",
        python_callable=check_task_completed,
        op_args=['ecs_operator_task'],
        mode="reschedule",
        poke_interval=120,
        timeout=60*60*2, #2 hours
    )    

    sns_pass=SnsPublishOperator(
        task_id="sns_pass",
        target_arn=success_sns_arn,
        message=success_sns_message, 
        subject=success_sns_message,
        retries=1
    )

    # 修复任务依赖
    run_ecs_task >> await_task_finish >> sns_pass

关键修改点说明

  • 任务依赖:添加run_ecs_task >> await_task_finish,确保Sensor在ECS任务启动完成后再开始轮询。
  • XCom提取逻辑:将xcom.get('task_arn')改为xcom['tasks'][0]['taskArn'],匹配boto3响应的正确键名。
  • 容错处理:增加XCom为空、任务查询失败的判断,返回False让Sensor继续等待,避免过早抛出异常。

内容的提问来源于stack exchange,提问作者Marco A. Falconi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 22:25:06