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

Airflow DAG完成后触发多次额外运行的原因及解决方法

问题描述

我在Airflow DAG中配置了一个将数据从BigQuery导出至GCS文件的任务(导出URI使用了通配符*),发现导出任务执行完毕后,会触发DAG的多次额外运行,且这些额外运行的首个任务都失败了。额外运行的次数和导出的文件数量一致;如果移除URI中的通配符,导出成功后只会触发一次额外运行。

Airflow树形视图:
Airflow树形视图展示

对应的DAG代码:

import datetime
import json
from airflow import utils
from airflow.models import DAG
from airflow.operators.dummy import DummyOperator
from airflow.contrib.operators.bigquery_operator import BigQueryOperator
from airflow.operators.python_operator import PythonOperator
from google.cloud import bigquery, pubsub_v1


def publish_message(**context):
    publisher = pubsub_v1.PublisherClient()

    topic_path = publisher.topic_path('a', 'a')

    data_str = json.dumps(context['task_instance'].xcom_pull(task_ids='set_params', key='request_params'))
    data = data_str.encode("utf-8")

    future = publisher.publish(topic_path, data)
    print(future.result())

    print(f"Published messages to {topic_path}.")


def set_params(params, **context):
    query_params = {}
    file_format = 'json'
    compression = 'gzip'
    params_fixed = json.loads(params.replace("\'", "\""))
    for sub in params_fixed['params']:
        if sub["param_name"] == 'file_format':
            file_format = sub["param_value"]
        elif sub["param_name"] == 'compression':
            compression = sub["param_value"]
        else:
            query_params[sub["param_name"]] = sub["param_value"]
    context['task_instance'].xcom_push(key='file_format', value=file_format)
    context['task_instance'].xcom_push(key='compression', value=compression)
    context['task_instance'].xcom_push(key='query_params', value=query_params)
    context['task_instance'].xcom_push(key='request_params', value=params_fixed)
    

with DAG(dag_id="export_data_to_gcs",
        start_date=datetime.datetime(2021, 1, 1),
        schedule_interval=None) as dag:

    start = DummyOperator(task_id="start")
    
    my_params = '{{dag_run.conf}}'
    
    set_params = PythonOperator(
        task_id='set_params',
        provide_context=True,
        python_callable=set_params,
        op_kwargs={'params': my_params},
        dag=dag,
    )  
    
    execute_sql = BigQueryOperator(
        task_id='execute_sql',
        sql='/SQL/sql_query.sql',
        bigquery_conn_id='bigquery_default',
        use_legacy_sql=False,
        dag=dag,
    )

    publish_message = PythonOperator(
        task_id='publish_message',
        provide_context=True,
        python_callable=publish_message,
        dag=dag,
    )
   
   
start >> set_params >> execute_sql >> publish_message
原因分析

这是因为使用通配符导出时,BigQuery会生成多个分片输出文件,每个文件完成后都会向Airflow发送作业完成回调通知。而BigQueryOperator的默认逻辑会把每个回调事件触发为一次新的DAG运行,但这些新运行没有携带原始DAG启动时的dag_run.conf参数,导致set_params任务在解析{{dag_run.conf}}时因参数缺失而失败。

如果不使用通配符,BigQuery仅生成单个文件,只会发送一次回调通知,因此只会触发一次额外运行。

解决方案

1. 使用专用导出算子替代BigQueryOperator

推荐使用BigQueryToGCSOperator,它是Airflow专为BigQuery导出到GCS场景设计的算子,不会触发多余的DAG运行:

from airflow.contrib.operators.bigquery_to_gcs import BigQueryToGCSOperator

# 替换原来的execute_sql任务
export_to_gcs = BigQueryToGCSOperator(
    task_id='export_to_gcs',
    source_project_dataset_table='your_project.your_dataset.your_table',
    destination_cloud_storage_uris=['gs://your_bucket/path/*'],
    export_format='JSON',
    compression='GZIP',
    bigquery_conn_id='bigquery_default',
    dag=dag
)

# 更新任务依赖
start >> set_params >> export_to_gcs >> publish_message

2. 为set_params任务添加参数容错处理

如果必须保留BigQueryOperator,可以在set_params函数中增加空值判断,避免因dag_run.conf缺失导致任务失败:

def set_params(params, **context):
    # 处理dag_run.conf为空的情况
    if not params or params == '{}':
        context['task_instance'].xcom_push(key='request_params', value={'params': []})
        return
    
    # 原有逻辑保留
    query_params = {}
    file_format = 'json'
    compression = 'gzip'
    params_fixed = json.loads(params.replace("\'", "\""))
    for sub in params_fixed['params']:
        if sub["param_name"] == 'file_format':
            file_format = sub["param_value"]
        elif sub["param_name"] == 'compression':
            compression = sub["param_value"]
        else:
            query_params[sub["param_name"]] = sub["param_value"]
    context['task_instance'].xcom_push(key='file_format', value=file_format)
    context['task_instance'].xcom_push(key='compression', value=compression)
    context['task_instance'].xcom_push(key='query_params', value=query_params)
    context['task_instance'].xcom_push(key='request_params', value=params_fixed)

3. 调整Airflow全局回调配置(谨慎使用)

可以在Airflow配置文件中修改bigquery_callback相关设置,阻止多余的DAG运行触发,但这种方式会影响所有BigQuery任务,需根据实际场景评估使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 18:25:25