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

Airflow初始化脚本自动执行问题排查及修正咨询

问题根源分析

你的脚本在Airflow初始化时就自动执行插入操作,核心问题出在两个地方:

  1. 全局作用域直接调用业务函数
    不管是原始脚本里的db_login()、insert_data(),还是更新后的load_etl(),这些函数调用写在了DAG定义之外的全局代码中。Airflow在解析DAG文件时,会从头到尾执行所有全局代码,所以这些函数会在DAG加载阶段就自动运行,完全不受任务调度逻辑控制。

  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 06:43:30