Airflow初始化脚本自动执行问题排查及修正咨询
问题根源分析
你的脚本在Airflow初始化时就自动执行插入操作,核心问题出在两个地方:
全局作用域直接调用业务函数
不管是原始脚本里的db_login()、insert_data(),还是更新后的load_etl(),这些函数调用写在了DAG定义之外的全局代码中。Airflow在解析DAG文件时,会从头到尾执行所有全局代码,所以这些函数会在DAG加载阶段就自动运行,完全不受任务调度逻辑控制。Operator使用错误
你用了BashOperator却又设置了python_callable参数——BashOperator是用来执行shell命令的,要运行Python函数应该用PythonOperator。而且就算用对了Operator,你给python_callable传的是load_etl()(带括号的函数调用),这会导致函数在DAG加载时就被执行,而非任务触发时才调用。
修复方案
以下是修正后的完整脚本,关键修改点已标注:
## Third party Library Imports import psycopg2 from airflow import DAG from airflow.operators.python import PythonOperator # 替换为PythonOperator from datetime import datetime, timedelta default_args = { 'owner': 'admin', 'start_date': datetime(2018, 5, 25), 'email': ['admin@mail.com'], 'email_on_failure': True, 'email_on_retry': True, 'retries': 1, 'retry_delay': timedelta(minutes=1), } dag = DAG('sample', default_args=default_args, catchup=False, schedule_interval="@once") def db_login(): try: db_con = psycopg2.connect( "dbname='db' user='user' password='password' host='host' port='5439' sslmode='require'" ) print('Connection success') return db_con except Exception as e: print(f"I am unable to connect to the database: {str(e)}") raise # 抛出异常让Airflow标记任务失败 def insert_data(db_con): cur = db_con.cursor() try: cur.execute("""insert into table_1 select id,name,status from table_2 limit 2;""") db_con.commit() print('ETL Task Complete: Inserting data into table_1') except Exception as e: db_con.rollback() # 出错时回滚事务 print(f"Insert failed: {str(e)}") raise finally: cur.close() db_con.close() # 确保连接最终关闭 def load_etl(): db_con = db_login() insert_data(db_con) # 移除全局作用域的load_etl()调用!避免DAG加载时自动执行 # 使用PythonOperator定义任务 t1 = PythonOperator( task_id='run_etl', python_callable=load_etl, # 传递函数名,不带括号! email_on_failure=True, email=['admin@mail.com'], dag=dag ) t1
关键修改说明
- 移除全局函数调用:删掉了原来的
load_etl(),彻底避免DAG加载时自动执行业务逻辑 - 替换为PythonOperator:用
PythonOperator替代BashOperator,这是Airflow中执行Python业务逻辑的标准方式 - python_callable传函数名:给
python_callable传递load_etl(不带括号),保证Airflow仅在任务触发时才调用该函数 - 优化数据库连接管理:去掉全局变量
db_con,改成函数返回连接并传递,避免资源泄漏;增加异常回滚逻辑,保证数据一致性;确保连接和游标最终关闭 - 移除无效参数:
PythonOperator不需要bash_command参数,直接删除即可
额外注意事项
- Airflow会定期解析DAG文件(默认每30秒一次),所以绝对不要在全局作用域编写任何有副作用的代码(比如数据库写入、文件修改等)
- 如果需要复用数据库连接,推荐使用Airflow UI的连接管理功能配置数据库连接,再用
PostgresHook获取连接,这比手动编写psycopg2连接更安全、更符合Airflow最佳实践
内容的提问来源于stack exchange,提问作者dark horse
相关产品推荐
相关产品推荐

