Airflow DAG启动时自动执行异常问题(版本2.2.5)
问题根源
你的代码中PythonOperator的python_callable参数直接执行了get_to_know_details(job_type,SOURCE_PATH)函数调用,而非传递函数本身。Airflow启动时会自动解析所有DAG文件,解析到这行代码时就会立即执行该函数,而非等到任务触发阶段才执行。
解决方法
需要将函数引用传给python_callable,通过op_kwargs传递参数,或用functools.partial绑定参数,两种方式任选其一:
方式一:使用op_kwargs传递参数
修改PythonOperator配置,仅传递函数名给python_callable,参数通过op_kwargs传入:
from airflow.operators.python import PythonOperator from datetime import datetime from airflow import DAG job_type='daily' SOURCE_PATH='/home/ubuntu/daily_data' def get_to_know_details(job_type,SOURCE_PATH): print("************************",job_type,SOURCE_PATH) with DAG( dag_id="transformer_daily_v1", is_paused_upon_creation=False, default_args=default_args, description="transformer to insert data", start_date=datetime(2022,9,20), schedule_interval='31 12 * * *', catchup=False ) as dag: task1=PythonOperator( task_id="dag_task_1", python_callable=get_to_know_details, # 仅传递函数名,不执行调用 op_kwargs={"job_type": job_type, "SOURCE_PATH": SOURCE_PATH} # 通过op_kwargs传递参数 )
方式二:使用functools.partial绑定参数
若习惯提前绑定参数,可使用partial包装函数:
from airflow.operators.python import PythonOperator from datetime import datetime from airflow import DAG from functools import partial job_type='daily' SOURCE_PATH='/home/ubuntu/daily_data' def get_to_know_details(job_type,SOURCE_PATH): print("************************",job_type,SOURCE_PATH) # 提前绑定参数 bound_func = partial(get_to_know_details, job_type=job_type, SOURCE_PATH=SOURCE_PATH) with DAG( dag_id="transformer_daily_v1", is_paused_upon_creation=False, default_args=default_args, description="transformer to insert data", start_date=datetime(2022,9,20), schedule_interval='31 12 * * *', catchup=False ) as dag: task1=PythonOperator( task_id="dag_task_1", python_callable=bound_func # 传递绑定好参数的函数 )
补充说明
Airflow解析DAG文件时会执行所有顶层代码,因此直接的函数调用会在解析阶段触发,而非任务运行阶段。只有传递函数引用给python_callable,Airflow才会在任务触发时执行该函数。
内容的提问来源于stack exchange,提问作者Ritik Gupta
相关产品推荐
相关产品推荐

