SQLAlchemy异步模式插入数据失败问题求助
异步写入数据库失败,但表能正常创建的问题
我尝试用异步模式向数据库写入数据,数据无法存储,但表可以正常创建。我怀疑问题出在create_async_session方法,但找不到具体原因,尝试改用固定会话也无法解决。
数据库相关代码
async def create_tables(self): async with self.async_engine.begin() as conn: await conn.run_sync(self.models.Base.metadata.drop_all) await conn.run_sync(self.models.Base.metadata.create_all) async def connect(self): conn_str: str = self.__create_connection_str() # ASYNC self.async_engine = create_async_engine(conn_str) self.session_factory = sessionmaker( self.async_engine, class_=AsyncSession, expire_on_commit=False, autocommit=False, autoflush=False ) def create_async_session(self): session = self.session_factory() try: yield session finally: session.close()
create_async_session返回:<generator object DatabaseHandler.create_async_session at 0x000001C148797840>
插入函数代码
async def create_artist(session: AsyncSession, artist: schemas.Artist): try: # map values db_artist = models.Artist( artist_id=artist.id, name=artist.name, href=artist.href, genres=artist.genres, popularity=artist.popularity, type=artist.type, uri=artist.uri, external_urls=artist.external_urls["spotify"], followers=artist.followers["total"] ) print(db_artist) # -> 返回正常数据 print(session) # -> 有返回值 await session.add(db_artist) await session.commit() await session.refresh(db_artist) return artist except Exception as err: err
主函数代码
async def collect_artists_in_background(): await db_client.connect() await db_client.create_tables() # -> 正常执行 session = db_client.create_async_session() # -> 正常执行 for _id in list_of_artists_ids: artist = await crud.get_artist_by_id(_id) # -> 正常执行 await CRUD.create_artist(session, schemas.Artist.parse_obj(artist)) # -> 执行失败
我尝试改用固定会话,仍然无法正常工作:
async with db_client.session_factory() as session: for _id in list_of_artists_ids: artist = await crud.get_artist_by_id(_id) await CRUD.create_artist(session, schemas.Artist.parse_obj(artist))
问题排查与修复方案
1. 生成器会话获取错误
create_async_session是生成器函数,直接赋值给session得到的是生成器对象,不是实际的AsyncSession实例。需要用next()提取会话:
session = next(db_client.create_async_session())
但更推荐直接使用session_factory的异步上下文管理器,这是SQLAlchemy异步会话的标准用法。
2. 异常处理缺失关键步骤
create_artist中的except块仅捕获异常但未处理,既不打印错误信息也不回滚会话,导致你无法定位问题。修改如下:
async def create_artist(session: AsyncSession, artist: schemas.Artist): try: # ... 原有映射代码 ... await session.add(db_artist) await session.commit() await session.refresh(db_artist) return artist except Exception as err: print(f"插入失败详情: {err}") await session.rollback() # 异常必须回滚,避免会话处于无效状态 raise # 重新抛出异常让上层感知错误
3. 检查异步引擎连接字符串
确保连接字符串使用了异步驱动:
- PostgreSQL:
postgresql+asyncpg://用户名:密码@主机/数据库名 - MySQL:
mysql+aiomysql://用户名:密码@主机/数据库名
如果误用同步驱动的URL(比如postgresql://而非postgresql+asyncpg://),异步引擎会无法正确执行写入操作。
4. 验证会话有效性
在调用create_artist前,确认session是AsyncSession实例:
print(type(session)) # 应该输出 <class 'sqlalchemy.ext.asyncio.session.AsyncSession'>
如果类型不对,说明会话创建逻辑有误。
内容的提问来源于stack exchange,提问作者Anna
相关产品推荐
相关产品推荐

