Apache Airflow数据加载报错:'Engine'对象无'cursor'属性
问题:Airflow批处理任务Load阶段报错:'Engine' object has no attribute 'cursor'
问题背景
使用Apache Airflow执行批处理任务,抽取(Extract)和转换(Transform)阶段运行正常,但加载(Load)阶段出现报错。已降级SQLAlchemy和Pandas版本,当前依赖版本:
- apache-airflow-core==3.0.4
- SQLAlchemy==1.4.54
- SQLAlchemy-JSONField==1.0.2
- SQLAlchemy-Utils==0.41.2
- pandas==1.5.3
Load任务代码
from airflow.decorators import dag, task from airflow.hooks.base import BaseHook from airflow.utils.log.logging_mixin import LoggingMixin from sqlalchemy import create_engine import pandas as pd from sqlalchemy_utils import create_database, database_exists import requests from bs4 import BeautifulSoup from datetime import datetime import time from datetime import timedelta import psycopg2 @task def load(df): try: connection_details = BaseHook.get_connection("my_postgres_conn") except Exception as e: print(f"❌ Could not retrieve Airflow connection: {e}") raise uid = connection_details.login pwd = connection_details.password host = connection_details.host port = connection_details.port schema = connection_details.schema postgres_url = f"postgresql+psycopg2://{uid}:{pwd}@{host}:{port}/{schema}" # Create the database if it doesn't exist if not database_exists(postgres_url): create_database(postgres_url) engine = create_engine(postgres_url) try: # Use the engine object directly with to_sql # This syntax is fully compatible with SQLAlchemy 1.4.54 df.to_sql("Bank_Cap", con=engine, if_exists="replace", index=False) print("✅ Data loaded successfully into PostgreSQL") except Exception as e: print(f"❌ Error loading data: {e}") raise # Save a CSV as backup df.to_csv("my_largest_bank.csv", index=False) print(f"Postgres URL: {postgres_url}") print(f"DataFrame to load:\n{df}") log.info(f"Postgres URL: {postgres_url}") log.info(f"DataFrame to load: {df}")
报错信息
[2025-08-17, 02:06:00] WARNING - Using Connection.get_connection_from_secrets from <code>airflow.models</code> is deprecated.Please use <code>from airflow.sdk import Connection</code> instead: category="DeprecationWarning": filename="/home/eziuche/.local/share/virtualenvs/airflow-project-kVhV5bz8/lib/python3.11/site-packages/airflow/models/connection.py": lineno=471: source="py.warnings" [2025-08-17, 02:06:00] INFO - Connection Retrieved 'my_postgres_conn': source="airflow.hooks.base" [2025-08-17, 02:06:00] WARNING - pandas only supports SQLAlchemy connectable (engine/connection) or database string URI or sqlite3 DBAPI2 connection. Other DBAPI2 objects are not tested. Please consider using SQLAlchemy.: category="UserWarning": [2025-08-17, 02:06:00] INFO - ❌ Error loading data: 'Engine' object has no attribute 'cursor': chan="stdout": source="task" [2025-08-17, 02:06:00] ERROR - Task failed with exception: source="task" AttributeError: 'Engine' object has no attribute 'cursor'
问题分析与解决方案
核心原因
pandas 1.5.x版本在处理SQLAlchemy Engine对象时,部分场景下会错误尝试调用DBAPI连接专属的cursor()方法(Engine对象本身没有该方法),通常是依赖冲突或连接使用方式导致的兼容性问题。
修复方案
方案1:使用Engine连接实例替代Engine对象
修改df.to_sql的调用逻辑,通过engine.connect()获取可执行cursor的连接对象:
try: with engine.connect() as conn: df.to_sql("Bank_Cap", con=conn, if_exists="replace", index=False) print("✅ Data loaded successfully into PostgreSQL") except Exception as e: print(f"❌ Error loading data: {e}") raise
方案2:直接传入连接字符串
跳过创建Engine的步骤,直接将构造好的postgres_url传给to_sql:
try: df.to_sql("Bank_Cap", con=postgres_url, if_exists="replace", index=False) print("✅ Data loaded successfully into PostgreSQL") except Exception as e: print(f"❌ Error loading data: {e}") raise
额外优化:消除Airflow连接的 deprecation warning
替换旧的连接获取方式为Airflow 3.x推荐的API:
# 替换原有的from airflow.hooks.base import BaseHook from airflow.sdk import Connection try: connection_details = Connection.get_connection_from_secrets("my_postgres_conn") except Exception as e: print(f"❌ Could not retrieve Airflow connection: {e}") raise
验证注意事项
- 确保已安装
psycopg2-binary(虚拟环境中优先使用该包,避免系统依赖问题) - 确认连接字符串中的
schema参数对应正确的PostgreSQL数据库名称,避免混淆概念
内容的提问来源于stack exchange,提问作者Nwaogu Eziuche
相关产品推荐
相关产品推荐

