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'
这个错误是多个问题叠加导致的:
- 函数名不匹配:你定义的PythonOperator用了
python_callable=data_warehouse_login,但实际你写的连接函数叫db_log,这会导致Airflow找不到可执行的函数,返回None,进而引发后续错误。 - XCom使用错误:你在
db_log里push的是字符串"dwh_connection",而不是实际的数据库连接对象——更关键的是,数据库连接对象不能通过XCom传递,因为XCom需要序列化对象,而数据库连接是不可序列化的资源对象。 - SQL语句无效:
insert into tbl_1 select limit 2是语法错误的SQL,缺少要选择的字段和来源表。 - 变量名错误:
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
相关产品推荐
相关产品推荐

