Airflow单任务失败后如何继续执行其他应用的任务集
Airflow+DBT多应用任务处理:实现任务失败不阻塞其他应用执行
问题场景
我们使用Airflow结合DBT处理多个带唯一ID的应用数据,每个应用需要依次执行3个Airflow任务完成数据处理。当前代码将所有任务串成一条全局串行链,只要其中任一任务失败,后续所有任务都会被标记为跳过。需要调整任务依赖逻辑,实现单个应用的任务失败后,其他应用的首个任务仍能正常执行。
原代码问题分析
原代码通过chain_tasks列表收集所有应用的所有任务,最后通过循环为每个任务设置下游任务,形成了全局串行依赖链:
应用A任务1 → 应用A任务2 → 应用A任务3 → 应用B任务1 → 应用B任务2 → ...
这种结构下,只要链中某一任务失败,后续所有任务都会被Airflow标记为跳过,无法实现其他应用任务的独立执行。
解决方案
为每个应用单独构建内部任务依赖(同一应用的3个任务依次执行),不同应用的任务链之间不建立依赖关系。这样所有应用的任务链会并行启动执行,单个应用的任务失败只会阻断该应用内部的后续任务,不会影响其他应用的任务流程。
修改后的代码
import os import sys sys.path.insert(0, os.path.abspath(os.path.dirname(__file__))) from airflow import DAG from airflow.models import Variable from airflow.contrib.operators.kubernetes_pod_operator import KubernetesPodOperator from utils.base_util import (default_args) from utils.token_util import (fetch_token) from utils.backend_util import get_applications dag_id = 'winback' START_DATE = Variable.get("AIRFLOW_START_DATE") BO_URL = Variable.get("URL") USER_NAME = Variable.get("AIRFLOW_USER_ID", default_var=os.environ.get("AIRFLOW_USER_ID")) PASSWORD = Variable.get("AIRFLOW_USER_PASSWORD", default_var=os.environ.get("AIRFLOW_USER_PASSWORD")) ENV = os.environ.get("ENVIRONMENT") AWS_ACCESS_KEY_ID = os.environ['AWS_ACCESS_KEY_ID'] AWS_SECRET_ACCESS_KEY = os.environ['AWS_SECRET_ACCESS_KEY'] AWS_DEFAULT_REGION = os.environ['AWS_DEFAULT_REGION'] TAG_DBT_REVERSE_EL = Variable.get("TAG_DBT_REVERSE_EL") TENANT = Variable.get("TENANT", default_var='SAAS') ORG = os.environ.get("ORGANIZATION_NAME") token = fetch_token(BO_URL, USER_NAME, PASSWORD) application_list = get_applications(BO_URL, token) if TENANT == 'SAAS': dag = DAG( dag_id, default_args=default_args, is_paused_upon_creation=False, schedule_interval=None, catchup=False, tags=["SEGMENTS", "DBT"]) # 移除全局chain_tasks列表,改为每个应用单独处理依赖 for application in application_list: application_id = str(application['application_id']) if str(application['is_seg_enabled']) == '1' and str(application['is_act_enabled']) == '1': suffix = application_id[application_id.find("-") + 1:] env_var = { 'S3_STAGING_DIR': f"s3://{os.environ.get('QUERY_LOGS_BUCKET')}/dbt/", 'REGION_NAME': os.environ.get("AWS_DEFAULT_REGION"), 'ENV': ENV, 'ORG': ORG, 'TENANT': TENANT, 'AWS_DEFAULT_REGION': AWS_DEFAULT_REGION, 'AWS_ACCESS_KEY_ID': AWS_ACCESS_KEY_ID, 'AWS_SECRET_ACCESS_KEY': AWS_SECRET_ACCESS_KEY, 'APPLICATION_ID': application_id, # 保留原有的其他环境变量 } winback = KubernetesPodOperator(namespace='etl', image=f'blotout/db-ana:{TAG_DBT_ANALYTICS}', cmds=["/usr/local/bin/dbt"], arguments=['run', '--models', 'winback'], env_vars=env_var, name="winback", configmaps=['awskey'], task_id=f"winback_{suffix}", get_logs=True, dag=dag, is_delete_operator_pod=True, ) winback_activation_stats = KubernetesPodOperator(namespace='etl', image=f'blotout/db-ana:{TAG_DBT_ANALYTICS}', cmds=["/usr/local/bin/dbt"], arguments=['run', '--models', 'winback_activation_stats'], env_vars=env_var, name="winback_activation_stats", configmaps=['awskey'], task_id=f"winback_activation_stats_{suffix}", get_logs=True, dag=dag, is_delete_operator_pod=True, ) winback_segments_sync = KubernetesPodOperator(namespace='etl', image=f'blotout/rev-el:{TAG_DBT_REVERSE_EL}', cmds=["python3"], arguments=['providers/activation/segments_init.py', 'winback', application_id], env_vars=env_var, name="winback_segments_sync", configmaps=['awskey'], task_id=f"winback_segments_sync_{suffix}", get_logs=True, dag=dag, is_delete_operator_pod=True, ) # 为当前应用的任务设置内部依赖:winback → winback_activation_stats → winback_segments_sync winback >> winback_activation_stats >> winback_segments_sync # 移除全局设置下游的循环 globals()[dag_id] = dag
关键修改点
- 移除了全局的
chain_tasks列表,不再将所有任务收集到同一列表中 - 在每个应用的循环内部,直接通过
>>操作符设置该应用三个任务的串行依赖 - 删除了最后为所有任务设置全局下游的循环逻辑
这样调整后,每个应用的三个任务会按顺序执行,不同应用的任务链之间相互独立,某一应用的任务失败后,其他应用的任务仍能正常启动执行。
内容的提问来源于stack exchange,提问作者azaveri7
相关产品推荐
相关产品推荐

