Airflow 2.2.5 DAG解析触发TaskTimeout超时问题咨询
问题现象
在Airflow Version 2.2.5/Composer 2.0.15环境中出现DAG解析失败报错,完全相同的业务代码在Airflow version2.2.3 /Composer Version 1.18.0环境中可正常运行。
核心报错为DagBag导入超时,具体报错栈如下:
Broken DAG: [/home/airflow/gcs/dags/test_dag.py] Traceback (most recent call last): File "/opt/python3.8/lib/python3.8/enum.py", line 256, in __new__ if canonical_member._value_ == enum_member._value_: File "/opt/python3.8/lib/python3.8/site-packages/airflow/utils/timeout.py", line 37, in handle_timeout raise AirflowTaskTimeout(self.error_message) airflow.exceptions.AirflowTaskTimeout: DagBag import timeout for /home/airflow/gcs/dags/test_dag.py after 30.0s. Please take a look at these docs to improve your DAG import time: * https://airflow.apache.org/docs/apache-airflow/2.2.5/best-practices.html#top-level-python-code * https://airflow.apache.org/docs/apache-airflow/2.2.5/best-practices.html#reducing-dag-complexity, PID: 1827
注:报错栈中enum模块相关的报错只是超时信号触发时代码刚好执行到该位置,不是enum模块本身存在故障。
现有项目结构
当前项目根目录main_folder下分三个子目录,职责如下:
- dags:存放所有主DAG文件
- tasks:存放PythonOperator对应执行函数、SQL查询等实际业务逻辑的
.py文件 - libs:存放通用公共功能实现的Python代码文件
提供的基础DAG结构示例代码如下:
# Import libraries and functions import datetime from airflow import models, DAG from airflow.contrib.operators import bigquery_operator, bigquery_to_gcs, bigquery_table_delete_operator from airflow.operators.python_operator import PythonOperator from airflow.operators.bash_operator import BashOperator ##from airflow.executors.sequential_executor import SequentialExecutor from airflow.utils.task_group import TaskGroup ## Import codes from tasks and libs folder from libs.compres_suppress.cot_suppress import * from libs.teams_plugin.teams_plugin import * from tasks.email_code.trigger_email import * # Set up Airflow DAG default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime.datetime(2020, 12, 15, 0), 'retries': 1, 'retry_delay': datetime.timedelta(minutes=1), 'on_failure_callback': trigger_email } DAG_ID = 'test_dag' # Check exscution date if "<some condition>" matches: run_date = <date in config file> else: run_date = datetime.datetime.now().strftime("%Y-%m-%d") run_date_day = datetime.datetime.now().isoweekday() dag = DAG( DAG_ID, default_args=default_args, catchup=False, max_active_runs=1, schedule_interval=SCHEDULE_INTERVAL ) next_dag_name = "next_dag1" if env == "prod": if run_date_day == 7: next_dag_name = "next_dag2" else: next_dag_name = "next_dag1" run_id = datetime.datetime.now().strftime("%Y-%m-%dT%H:%M:%S") # Define Airflow DAG with dag: team_notify_task = MSTeamsWebhookOperator( task_id='teams_notifi_start_task', http_conn_id='http_conn_id', message=f"DAG has started <br />" f"<strong> DAG ID:</strong> {DAG_ID}.<br />", theme_color="00FF00", button_text="My button", dag=dag) task1_bq = bigquery_operator.BigQueryOperator( task_id='task1', sql=task1_query( table1="table1", start_date=start_date), use_legacy_sql=False, destination_dataset_table="destination_tbl_name", write_disposition='WRITE_TRUNCATE' ) ##### Base Skeletons ##### with TaskGroup("taskgroup_lbl", tooltip="taskgroup_sample") as task_grp: tg_process(args=default_args,run_date=run_date) if run_mode == "<env_name>" and next_dag != "": next_dag_trigg = BashOperator( task_id=f'trigger_{next_dag_name}', bash_command="gcloud composer environments run " + <env> + "-cust_comp --location us-east1 dags trigger -- " + next_dag_name + " --run-id='trigger_ '" ) task_grp >> next_dag_trigger team_notify_task >> task1_bq >> task_grp
根因定位
- 版本机制差异
- Airflow 2.2.3版本存在DagBag导入超时计时逻辑bug,实际导入耗时超过30s也不会触发超时报错;2.2.4及之后版本修复了该bug,超时判定逻辑恢复正常。
- Composer 1.x的DAG目录为本地磁盘挂载,文件读取、模块导入IO延迟极低;Composer 2.x默认通过GCSfuse挂载DAG存储桶,跨目录导入自定义模块的IO开销比1.x高数倍,相同代码的导入耗时会明显变长。
- 代码写法问题
- 大量使用
import *语法导入libs、tasks下的模块,会执行对应模块所有顶层代码,包括DAG解析阶段不需要的业务逻辑、第三方库初始化逻辑,无意义拉长导入时间。 - DAG顶层直接执行动态时间计算、配置读取、
task1_querySQL生成、tg_processTaskGroup生成逻辑,如果这些逻辑涉及读文件、拉取外部元数据、加载大段文本,全部会计入DAG导入耗时。 - 存在大量无效导入:示例代码中导入了
bigquery_to_gcs、bigquery_table_delete_operator、PythonOperator等未实际使用的对象,额外增加加载开销。
- 大量使用
- 隐性放大因素
如果libs、tasks目录下的模块存在顶层执行的IO逻辑、重型计算,或者导入了pandas、numpy、全量云服务客户端这类加载耗时较长的第三方库,在GCSfuse的IO延迟加持下,很容易突破30s的默认超时阈值。
修复方案
- 替换所有
import *的导入写法,仅导入DAG文件中实际用到的对象,无关依赖全部删除。 - 优化顶层代码逻辑:所有不需要在DAG解析阶段执行的逻辑,全部挪到任务对应的执行函数中做延迟加载,不要在DAG顶层执行任何文件读取、外部接口调用、重型计算操作。
- 调整超时配置兜底:在Composer的Airflow配置覆盖项中,将
[core] dagbag_import_timeout参数从默认30s调整为匹配项目代码体量的值(推荐60s-120s,不要设置过高避免异常DAG阻塞整个解析流程)。 - 排查所有自定义模块的顶层代码,禁止在模块导入阶段执行任何业务逻辑,所有功能全部封装为函数,仅在被主动调用时执行。
内容的提问来源于stack exchange,提问作者codninja0908
相关产品推荐
相关产品推荐

