Airflow执行Redshift SQL报错:AttributeError: 'NoneType'无execute属性
解决Airflow中Redshift执行SQL时的AttributeError问题
首先,咱们来拆解你遇到的AttributeError: 'NoneType' object has no attribute 'execute'错误,根源出在这几个关键问题上:
问题分析
- 函数名不匹配:你定义的数据库连接函数是
db_log,但在t1任务里指定的python_callable=data_warehouse_login,这会导致Airflow找不到对应的执行函数,大概率这个连接任务实际执行失败了,后续xcom_pull自然拿到的是None。 - XCom传递的是无效值:你在
db_log里用xcom_push存的是字符串"dwh_connection",就算这个值能正常传递,字符串也没有execute方法;而且函数最后return (dwh_connection)里的dwh_connection根本没定义,这会直接引发NameError,导致任务出错,最终XCom里的值变成None。 - 数据库连接不能通过XCom传递:就算你想传递实际的
psycopg2连接对象,XCom是通过序列化存储的,数据库连接这类带状态的对象无法被序列化,所以这种方式根本行不通。 - SQL语法错误:你的
insert into tbl_1 select limit 2 ;缺少要查询的来源表,比如应该是insert into tbl_1 select * from your_source_table limit 2,这后续执行也会报错。
修正方案(推荐使用Airflow官方Hook)
Airflow提供了专门的RedshiftHook来管理Redshift连接,比自己手动用psycopg2更可靠,也避免了连接传递的问题:
步骤1:在Airflow UI中配置Redshift连接
先在Airflow的Admin -> Connections里添加Redshift连接,填写对应的dbname、user、password、host、port等信息,连接ID设为redshift_default(或者自定义ID)。
步骤2:修改代码使用RedshiftHook
## Third party Library Imports import pandas as pd import airflow from airflow import DAG from airflow.operators.python_operator import PythonOperator from airflow.providers.amazon.aws.hooks.redshift import RedshiftHook from datetime import datetime, timedelta # Following are defaults which can be overridden later on default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2018, 5, 29, 12), 'email': ['airflow@airflow.com'] } dag = DAG('sample1', default_args=default_args) def insert_data(**kwargs): # 使用RedshiftHook获取连接 redshift_hook = RedshiftHook(redshift_conn_id='redshift_default') conn = redshift_hook.get_conn() cur = conn.cursor() # 修正SQL语法,替换成你的实际来源表名 try: cur.execute("""insert into tbl_1 select * from your_source_table limit 2""") conn.commit() # 记得提交事务 print("数据插入成功") except Exception as e: conn.rollback() # 出错时回滚事务 print(f"执行出错: {str(e)}") raise e finally: # 关闭游标和连接,释放资源 cur.close() conn.close() t2 = PythonOperator( task_id='insert_into_redshift', python_callable=insert_data, provide_context=True, dag=dag ) t2
如果你坚持手动管理连接(不推荐)
如果一定要用psycopg2手动处理,不要用XCom传递连接,而是在insert_data函数里重新创建连接,同时修正之前的错误:
## Third party Library Imports import pandas as pd import psycopg2 import airflow from airflow import DAG from airflow.operators.python_operator import PythonOperator from datetime import datetime, timedelta default_args = { 'owner': 'airflow', 'depends_on_past': False, 'start_date': datetime(2018, 5, 29, 12), 'email': ['airflow@airflow.com'] } dag = DAG('sample1', default_args=default_args) def insert_data(**kwargs): db_con = None cur = None try: # 直接在当前函数里创建数据库连接 db_con = psycopg2.connect( dbname='name', user='user', password='pass', host='host', port='5439' ) cur = db_con.cursor() # 修正SQL语法 cur.execute("""insert into tbl_1 select * from your_source_table limit 2""") db_con.commit() print("数据插入成功") except Exception as e: if db_con: db_con.rollback() print(f"执行出错: {str(e)}") raise e finally: # 确保游标和连接被关闭 if cur: cur.close() if db_con: db_con.close() t2 = PythonOperator( task_id='insert_into_redshift', python_callable=insert_data, provide_context=True, dag=dag ) t2
额外注意点
- 永远记得在数据库操作后提交事务(
commit()),出错时回滚(rollback()),避免数据不一致。 - 用完连接和游标后一定要关闭,防止资源泄漏。
- Airflow的任务可能运行在不同的Worker节点上,全局变量
db_con无法跨节点共享,所以不要用全局变量传递连接。
内容的提问来源于stack exchange,提问作者dark horse
相关产品推荐
相关产品推荐

