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

Airflow 2.3.0 TaskFlow API:上游任务失败触发清理任务的问题

解决方案

核心思路

让关闭EC2实例的TaskX仅依赖启动实例的TaskA输出,通过设置触发规则确保中间任务任一失败时触发执行,同时避免依赖中间任务的输出导致XCom获取失败的问题。

具体实现步骤

  1. 定义启动实例任务:TaskA启动EC2实例并返回实例ID,TaskFlow会自动将返回值存入XCom。
  2. 定义中间业务任务:所有中间任务(B-W)接收TaskA返回的实例ID作为参数,执行各自业务逻辑。
  3. 定义关闭实例任务:TaskX仅接收TaskA的实例ID作为参数,设置触发规则为TriggerRule.ONE_FAILED(仅中间任务失败时执行)或TriggerRule.ALL_DONE(无论中间任务成功/失败都执行),实现关闭EC2实例的逻辑。
  4. 设置依赖关系:构建TaskA >> 所有中间任务 >> TaskX的依赖链,确保中间任务完成后触发TaskX。

代码示例

from airflow import DAG
from airflow.decorators import task
from airflow.utils.trigger_rule import TriggerRule
from datetime import datetime
import boto3

# 启动EC2实例任务
@task
def start_ec2_instance():
    ec2_client = boto3.client("ec2", region_name="us-east-1")
    resp = ec2_client.run_instances(
        ImageId="ami-0c55b159cbfafe1f0",
        InstanceType="t2.micro",
        MinCount=1,
        MaxCount=1
    )
    return resp["Instances"][0]["InstanceId"]

# 中间业务任务示例(TaskB)
@task
def task_b(instance_id):
    print(f"Processing instance {instance_id} in TaskB")
    # 可手动触发失败测试:raise ValueError("TaskB failed manually")

# 中间业务任务示例(TaskC)
@task
def task_c(instance_id):
    print(f"Processing instance {instance_id} in TaskC")

# 关闭EC2实例任务
@task(trigger_rule=TriggerRule.ONE_FAILED)
def stop_ec2_instance(instance_id):
    ec2_client = boto3.client("ec2", region_name="us-east-1")
    ec2_client.terminate_instances(InstanceIds=[instance_id])
    print(f"Terminated EC2 instance {instance_id}")

# 定义DAG
with DAG(
    dag_id="ec2_workflow",
    start_date=datetime(2023, 1, 1),
    schedule_interval=None,
    catchup=False
) as dag:
    instance_id = start_ec2_instance()
    
    # 实例化所有中间任务
    task_b_run = task_b(instance_id)
    task_c_run = task_c(instance_id)
    # ... 依次实例化TaskD到TaskW
    
    # 设置依赖链
    instance_id >> [task_b_run, task_c_run] >> stop_ec2_instance(instance_id)

关键说明

  • 触发规则选择:
    • TriggerRule.ONE_FAILED:仅当中间任务中有任意一个失败时,TaskX才执行。
    • TriggerRule.ALL_DONE:无论中间任务成功或失败,只要所有中间任务完成(包括失败状态),TaskX都会执行(适合必须确保实例最终被关闭的场景)。
  • 避免XCom错误:TaskX仅依赖TaskA的实例ID,不读取任何中间任务的输出,因此不会因中间任务失败导致XCom未生成而报错。
  • 依赖语法正确性:使用列表语法[task_b_run, task_c_run]批量设置中间任务依赖,再通过>>连接到TaskX,完全符合Airflow 2.3.0 TaskFlow API的语法规范,不会触发TypeError。

内容的提问来源于stack exchange,提问作者inspectorG4dget

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 12:35:36