在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
相关产品推荐
相关产品推荐

