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版本。
核心问题分析
- 任务依赖缺失:原DAG中
run_ecs_task与await_task_finish未设置依赖,导致Sensor与ECSOperator并行执行,此时ECSOperator尚未推送XCom,故拉取结果为空。 - Task ARN提取路径错误:boto3的
run_task响应中,Task ARN的键为taskArn(驼峰命名),而非代码中使用的task_arn。 - 空值未处理:未对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
相关产品推荐
相关产品推荐

