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

如何基于SQLAlchemy异步使用pd.to_sql()写入PostgreSQL?

异步写入Pandas DataFrame到PostgreSQL的正确方案

问题背景

现有代码使用同步create_engine可正常运行,但改用create_async_engine后抛出'AsyncEngine object has no cursor'错误,需实现简洁、安全的异步写入方案。

原代码:

from datetime import datetime
import pandas as pd 
import asyncio

x = 0

async def write_sql(engine):
    global x 
    x += 1
    print(f"Running {x}") 
    d = {'col1': [1, 2, 3, 4], 'col2': [3, 4, 5, 6]}
    df = pd.DataFrame(data=d)   
    df.to_sql(con=engine, name="test_data", if_exists="replace", index=False)
    await asyncio.sleep(1)


async def main():    
    engine = create_engine('postgresql+asyncpg://user:pass@localhost:5432/postgres')
    t1 = datetime.now() 
    await asyncio.gather(write_sql(engine),write_sql(engine),write_sql(engine))
    t2 = datetime.now()
    delta = t2 - t1
    print(f"Took {delta} seconds")


if __name__=='__main__': 
    asyncio.run(main()) 
    print("Finished.")

错误原因

Pandas的df.to_sql()是同步方法,依赖传统同步DBAPI的连接/游标接口。而create_async_engine生成的AsyncEngine是为异步SQLAlchemy设计的,不具备同步方法所需的cursor属性,直接传入会触发报错。

正确解决方案

利用SQLAlchemy异步引擎的run_sync()方法,将同步的to_sql操作包装到异步上下文执行,既复用异步连接池优势,又兼容Pandas同步API。同时优化代码结构,提升安全性。

修改后的完整代码:

from datetime import datetime
import pandas as pd
import asyncio
from sqlalchemy.ext.asyncio import create_async_engine

async def write_sql(engine, task_id):
    print(f"Running task {task_id}")
    d = {'col1': [1, 2, 3, 4], 'col2': [3, 4, 5, 6]}
    df = pd.DataFrame(data=d)
    
    # 用异步连接的run_sync适配同步to_sql操作
    async with engine.begin() as conn:
        await conn.run_sync(
            lambda sync_conn: df.to_sql(
                name="test_data",
                con=sync_conn,
                if_exists="replace",
                index=False
            )
        )
    await asyncio.sleep(1)

async def main():
    # 建议从环境变量读取敏感信息,避免硬编码
    # import os
    # db_url = os.getenv("DATABASE_URL", "postgresql+asyncpg://user:pass@localhost:5432/postgres")
    db_url = "postgresql+asyncpg://user:pass@localhost:5432/postgres"
    engine = create_async_engine(db_url, echo=False)
    
    t1 = datetime.now()
    # 传递任务ID替代全局变量,避免多任务竞争
    await asyncio.gather(
        write_sql(engine, 1),
        write_sql(engine, 2),
        write_sql(engine, 3)
    )
    t2 = datetime.now()
    delta = t2 - t1
    print(f"Took {delta} seconds")
    
    # 释放异步引擎连接池资源
    await engine.dispose()

if __name__ == '__main__':
    asyncio.run(main())
    print("Finished.")

关键优化说明

  • 异步兼容:通过conn.run_sync()将同步写入操作适配到异步上下文,利用异步连接池管理连接。
  • 安全性:避免硬编码数据库密码,推荐用环境变量读取敏感配置。
  • 代码健壮性:用任务ID替代全局变量,消除多任务下的变量竞争风险。
  • 资源管理:用async with engine.begin()自动管理事务,任务结束后调用engine.dispose()释放资源。

额外注意事项

  • 写入大量数据时,添加chunksize参数分块处理,避免内存占用过高。
  • 开启echo=True可打印SQL日志,便于调试生产环境问题。

内容的提问来源于stack exchange,提问作者splotsh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 13:25:39