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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 12:12:27