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

Airflow执行Python函数脚本报错求助:NameError与AttributeError

让我们一步步解决你遇到的这两个Airflow问题,先拆解错误原因,再给出可运行的修正方案:

第一个错误:NameError: name 'task_instance' is not defined

这个错误出现在你最初的insert_data函数里,原因很直接:你没有像db_log函数那样通过kwargs['task_instance']来获取Airflow的任务实例对象,直接就使用了task_instance变量,Python自然找不到这个未定义的变量。

第二个错误:AttributeError: 'NoneType' object has no attribute 'execute'

这个错误是多个问题叠加导致的:

  1. 函数名不匹配:你定义的PythonOperator用了python_callable=data_warehouse_login,但实际你写的连接函数叫db_log,这会导致Airflow找不到可执行的函数,返回None,进而引发后续错误。
  2. XCom使用错误:你在db_log里push的是字符串"dwh_connection",而不是实际的数据库连接对象——更关键的是,数据库连接对象不能通过XCom传递,因为XCom需要序列化对象,而数据库连接是不可序列化的资源对象。
  3. SQL语句无效:insert into tbl_1 select limit 2是语法错误的SQL,缺少要选择的字段和来源表。
  4. 变量名错误:db_log函数最后返回的dwh_connection根本没定义,应该是你创建的db_con。

修正后的完整代码(遵循Airflow最佳实践)

我们不用XCom传递连接对象,而是每个任务独立创建连接(用Airflow内置的连接管理存储敏感信息,避免硬编码):

首先,先在Airflow UI的Admin > Connections里创建一个名为my_postgres_db的连接,填入你的数据库信息(dbname、user、password、host、port等)。

然后用下面的代码替换你的脚本:

## Third party Library Imports
import psycopg2
from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from airflow.hooks.base_hook import BaseHook
from datetime import datetime, timedelta

# Default arguments for the DAG
default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2018, 5, 29, 12),
    'email': ['airflow@airflow.com'],
    'email_on_failure': False,
    'email_on_retry': False,
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

# Initialize the DAG
dag = DAG('sample1', default_args=default_args, catchup=False, schedule_interval="@once")

# Helper function to get DB connection from Airflow Connections
def get_db_connection():
    # Pull connection config from Airflow's built-in connection manager
    conn = BaseHook.get_connection('my_postgres_db')
    db_con = psycopg2.connect(
        dbname=conn.schema,
        user=conn.login,
        password=conn.password,
        host=conn.host,
        port=conn.port
    )
    return db_con

def db_log(**kwargs):
    try:
        # Create connection and verify it works
        db_con = get_db_connection()
        print('Database connection task completed successfully')
        # Push a success flag instead of connection object
        kwargs['task_instance'].xcom_push(key="db_connected", value=True)
        db_con.close()  # Always close connections when done
    except Exception as e:
        print(f"Failed to connect to database: {str(e)}")
        kwargs['task_instance'].xcom_push(key="db_connected", value=False)
        raise  # Raise exception to mark task as failed in Airflow

def insert_data(**kwargs):
    # Check if previous connection task succeeded
    is_connected = kwargs['task_instance'].xcom_pull(key="db_connected", task_ids='DWH_Connect')
    if not is_connected:
        raise Exception("Aborting: Database connection failed in previous task")
    
    try:
        db_con = get_db_connection()
        cur = db_con.cursor()
        # Corrected SQL statement (adjust fields/tables to your actual needs)
        cur.execute("""insert into tbl_1 select id, bill_no, status from tbl_2 limit 2;""")
        db_con.commit()  # Commit the transaction
        print("Data inserted successfully")
        cur.close()
        db_con.close()
    except Exception as e:
        db_con.rollback()  # Rollback on failure
        print(f"Insert operation failed: {str(e)}")
        raise  # Raise exception to mark task as failed

# Define tasks
t1 = PythonOperator(
    task_id='DWH_Connect',
    python_callable=db_log,
    provide_context=True,
    dag=dag
)

t2 = PythonOperator(
    task_id='Insert_Data',
    python_callable=insert_data,
    provide_context=True,
    dag=dag
)

# Set task dependency
t1 >> t2

关键修改点说明

  • 使用Airflow连接管理:避免硬编码敏感的数据库凭证,更安全也便于维护。
  • 不传递连接对象:每个任务独立创建/关闭连接,避免序列化问题,也符合Airflow的任务隔离原则。
  • 事务管理:添加了commit和rollback,确保数据操作的一致性。
  • 错误处理:抛出异常让Airflow正确标记任务失败,而不是只打印错误信息。
  • 修正语法错误:修复了函数名不匹配、SQL语句无效等低级错误。

内容的提问来源于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 08:13:14