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

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

关键修改点

  1. 移除了全局的chain_tasks列表,不再将所有任务收集到同一列表中
  2. 在每个应用的循环内部,直接通过>>操作符设置该应用三个任务的串行依赖
  3. 删除了最后为所有任务设置全局下游的循环逻辑

这样调整后,每个应用的三个任务会按顺序执行,不同应用的任务链之间相互独立,某一应用的任务失败后,其他应用的任务仍能正常启动执行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 06:12:08