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

在AirFlow中使用df.to_sql写入PostgreSQL失败求助

AirFlow中pandas.to_sql迁移MSSQL数据到PostgreSQL报错解决

问题场景

运行两步AirFlow任务:

  • 从MSSQL服务器拉取数据转为DataFrame
  • 将DataFrame数据写入PostgreSQL数据库

使用MsSqlHook、PostgresHook管理连接,数据拉取正常,但.to_sql写入环节持续报错,此前用SQLAlchemy实现过同类操作但本次失败。

原代码

from airflow import DAG
from airflow.decorators import task
from datetime import timedelta
from airflow.providers.microsoft.mssql.hooks.mssql import MsSqlHook
from airflow.providers.postgres.hooks.postgres import PostgresHook

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': '2024-01-01',
    'email_on_failure': False,
    'email_on_retry': False,
    'retry_delay': timedelta(minutes=5),
}
    
def mssql_to_postgres_transfer(mssql_hook, postgres_hook, sql, table_name, dagrun_id):
    
    df = mssql_hook.get_pandas_df(sql=sql)
    df['dagrun_id'] = dagrun_id
    
    df.to_sql(
        name=table_name,
        con=postgres_hook.get_sqlalchemy_engine()
        if_exists='replace',
        )
    
def do_the_job():
    mssql_hook = MsSqlHook(mssql_conn_id='mssql_conn_id')
    postgres_hook = PostgresHook(postgres_conn_id='postgres_conn_id')
    
    qry = "SELECT * FROM tbl_dummy"
    
    mssql_to_postgres_transfer(
        mssql_hook=mssql_hook, 
        postgres_hook=postgres_hook, 
        sql=qry, 
        table_name='tbl_dummy_local', 
        dagrun_id=1
        )
    
with DAG(
    'my_dag',
    default_args=default_args,
    catchup=False,
    description='fetch and store',
    max_active_runs=1,
    schedule_interval='5 * * * *'
) as dag:
    
    @task
    def fetch_and_store_task():
        do_the_job()
    
    fetch_and_store_task()

报错1:AttributeError: 'Engine' object has no attribute 'cursor'

[2024-02-07T22:25:08.665+0000] {base.py:83} INFO - Using connection ID 'postgres_conn_id' for task execution.
Traceback (most recent call last):
  File "<stdin>", line 1, in <module>
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/util/_decorators.py", line 333, in wrapper
    return func(*args, **kwargs)
           ^^^^^^^^^^^^^^^^^^^^^
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/core/generic.py", line 3081, in to_sql
    return sql.to_sql(
           ^^^^^^^^^^^
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/io/sql.py", line 842, in to_sql
    return pandas_sql.to_sql(
           ^^^^^^^^^^^^^^^^^^
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/io/sql.py", line 2851, in to_sql
    table.create()
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/io/sql.py", line 984, in create
    if self.exists():
       ^^^^^^^^^^^^^
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/io/sql.py", line 970, in exists
    return self.pd_sql.has_table(self.name, self.schema)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/io/sql.py", line 2866, in has_table
    return len(self.execute(query, [name]).fetchall()) > 0
               ^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/io/sql.py", line 2673, in execute
    cur = self.con.cursor()
          ^^^^^^^^^^^^^^^
AttributeError: 'Engine' object has no attribute 'cursor'

尝试SQLAlchemy raw_connection后的报错2

尝试代码:

eng = postgres_hook.get_sqlalchemy_engine()
df.to_sql(name='tbl_dummy_local', con=eng.raw_connection(), if_exists='replace')

报错信息:

Traceback (most recent call last):
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/io/sql.py", line 2675, in execute
    cur.execute(sql, *args)
psycopg2.errors.UndefinedTable: relation "sqlite_master" does not exist
LINE 5:             sqlite_master
                    ^


The above exception was the direct cause of the following exception:

Traceback (most recent call last):
  File "<stdin>", line 1, in <module>
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/util/_decorators.py", line 333, in wrapper
    return func(*args, **kwargs)
           ^^^^^^^^^^^^^^^^^^^^^
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/core/generic.py", line 3081, in to_sql
    return sql.to_sql(
           ^^^^^^^^^^^
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/io/sql.py", line 842, in to_sql
    return pandas_sql.to_sql(
           ^^^^^^^^^^^^^^^^^^
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/io/sql.py", line 2851, in to_sql
    table.create()
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/io/sql.py", line 984, in create
    if self.exists():
       ^^^^^^^^^^^^^
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/io/sql.py", line 970, in exists
    return self.pd_sql.has_table(self.name, self.schema)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/io/sql.py", line 2866, in has_table
    return len(self.execute(query, [name]).fetchall()) > 0
               ^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/home/airflow/.local/lib/python3.11/site-packages/pandas/io/sql.py", line 2687, in execute
    raise ex from exc
pandas.errors.DatabaseError: Execution failed on sql '
        SELECT
            name
        FROM
            sqlite_master
        WHERE
            type IN ('table', 'view')
            AND name=?;
        ': relation "sqlite_master" does not exist
LINE 5:             sqlite_master
                    ^

验证信息

两个数据库连接正常:

>>> postgres_hook.test_connection()
[2024-02-07T22:25:51.755+0000] {base.py:83} INFO - Using connection ID 'aircerv' for task execution.
[2024-02-07T22:25:51.765+0000] {sql.py:450} INFO - Running statement: select 1, parameters: None
[2024-02-07T22:25:51.766+0000] {sql.py:459} INFO - Rows affected: 1
(True, 'Connection successfully tested')

DataFrame包含有效数据:

>>> df.shape
(366, 77)

解决方案

1. 修复代码语法错误

原代码中df.to_sql的参数存在语法问题:con参数后未加逗号,导致if_exists被错误解析为con参数的一部分。修正后:

df.to_sql(
    name=table_name,
    con=postgres_hook.get_sqlalchemy_engine(),  # 添加逗号分隔参数
    if_exists='replace',
)

2. 正确使用SQLAlchemy引擎

pandas的.to_sql对SQLAlchemy Engine对象的支持更完善,无需手动调用raw_connection()。修改后的核心函数:

def mssql_to_postgres_transfer(mssql_hook, postgres_hook, sql, table_name, dagrun_id):
    df = mssql_hook.get_pandas_df(sql=sql)
    df['dagrun_id'] = dagrun_id
    
    engine = postgres_hook.get_sqlalchemy_engine()
    df.to_sql(
        name=table_name,
        con=engine,
        if_exists='replace',
        index=False,  # 避免将DataFrame索引写入数据库表
        chunksize=1000  # 数据量大时分批写入,防止内存溢出
    )

3. 排查数据库连接配置

确保AirFlow中PostgreSQL连接的配置正确:

  • 连接类型选择Postgres
  • 填写正确的主机、端口、数据库名、用户名、密码
  • 若需指定schema,在.to_sql中添加schema参数(如schema='public')

4. 高性能替代方案(可选)

如果数据量较大,推荐使用PostgreSQL的COPY命令批量导入,性能远优于.to_sql:

from io import StringIO

def mssql_to_postgres_transfer(mssql_hook, postgres_hook, sql, table_name, dagrun_id):
    df = mssql_hook.get_pandas_df(sql=sql)
    df['dagrun_id'] = dagrun_id
    
    # 将DataFrame转为Tab分隔的CSV流
    buffer = StringIO()
    df.to_csv(buffer, index=False, header=False, sep='\t')
    buffer.seek(0)
    
    # 使用copy_expert执行批量导入
    copy_sql = f"""
        COPY {table_name} FROM stdin WITH CSV DELIMITER '\t' NULL ''
    """
    postgres_hook.copy_expert(sql=copy_sql, file=buffer)

内容的提问来源于stack exchange,提问作者Leonardo Checon Dantas

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 08:09:51