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

如何在Airflow中检测空文件时终止DAG并跳过后续处理?

问题

我正在开发一个Airflow DAG,仅当文件非空时才对其执行特定处理任务。理想的工作流应先检查文件是否有内容,若文件为空,则DAG需停止执行并跳过所有与该文件相关的后续处理。

以下是我的Airflow DAG简化结构:

from google.cloud import storage
from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from datetime import datetime

def check_file_not_empty():
    client = storage.Client()
    bucket = client.get_bucket(src_bucket_name)
    blob = bucket.get_blob(blob_name)
    if blob.size == 0:
        raise Exception(f"The file {blob_name} in bucket {src_bucket_name} is empty")

def process_file():
    # Code to process the file

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2024, 3, 21),
    'retries': 1,
}

dag = DAG('file_processing_dag', default_args=default_args, schedule_interval='@daily')

check_file_task = PythonOperator(
    task_id='check_file_not_empty',
    python_callable=check_file_not_empty,
    dag=dag,
)

process_file_task = PythonOperator(
    task_id='process_file',
    python_callable=process_file,
    dag=dag,
)

check_file_task >> process_file_task

目前抛出异常仅会触发重试,这并非我需要的效果。我希望实现快速失败,请问应调用Airflow的哪些内部选项来终止执行?


解决方案

1. 取消检查任务的重试机制

默认default_args里的retries=1会让失败任务自动重试,要实现快速失败,需要给检查任务单独设置retries=0,覆盖全局重试配置,这样任务失败后直接终止,不会触发重试。

2. 选择合适的异常类型控制下游任务状态

根据业务场景,有两种可选方案:

方案A:标记为异常失败(文件空属于错误场景)

保持抛出普通Exception,但给检查任务设置retries=0,此时检查任务会直接失败,下游任务会被标记为upstream_failed,整个DAG运行状态为失败。

方案B:自动跳过下游任务(文件空属于预期场景)

如果文件为空是业务允许的正常情况(比如无数据的日期),可以使用Airflow内置的AirflowSkipException,抛出该异常后,检查任务会被标记为“跳过”,所有依赖它的下游任务也会自动跳过,整个DAG运行状态为成功。

需要先导入该异常:

from airflow.exceptions import AirflowSkipException

3. 完整修改后的代码

方案B(推荐用于允许空文件的场景)

from google.cloud import storage
from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from airflow.exceptions import AirflowSkipException
from datetime import datetime

def check_file_not_empty():
    client = storage.Client()
    bucket = client.get_bucket(src_bucket_name)
    blob = bucket.get_blob(blob_name)
    if blob.size == 0:
        raise AirflowSkipException(f"The file {blob_name} in bucket {src_bucket_name} is empty, skipping downstream tasks")

def process_file():
    # 处理文件的代码
    pass

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2024, 3, 21),
    'retries': 1,
}

dag = DAG('file_processing_dag', default_args=default_args, schedule_interval='@daily')

check_file_task = PythonOperator(
    task_id='check_file_not_empty',
    python_callable=check_file_not_empty,
    dag=dag,
    retries=0  # 禁用重试,快速终止/跳过
)

process_file_task = PythonOperator(
    task_id='process_file',
    python_callable=process_file,
    dag=dag,
)

check_file_task >> process_file_task

方案A(用于文件空属于错误的场景)

只需将check_file_not_empty函数中的AirflowSkipException换回普通Exception即可,其余配置不变。

两种方案对比

方案检查任务状态下游任务状态DAG整体状态适用场景
方案A失败上游失败(upstream_failed)失败文件为空属于异常错误,需触发告警
方案B跳过跳过成功文件为空是预期情况,无需告警

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 16:30:32