Airflow DAG完成后触发多次额外运行的原因及解决方法
问题描述
我在Airflow DAG中配置了一个将数据从BigQuery导出至GCS文件的任务(导出URI使用了通配符*),发现导出任务执行完毕后,会触发DAG的多次额外运行,且这些额外运行的首个任务都失败了。额外运行的次数和导出的文件数量一致;如果移除URI中的通配符,导出成功后只会触发一次额外运行。
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
相关产品推荐
相关产品推荐

