如何用SQLAlchemy与AsyncSession正确实现get_all_by_chunk异步生成器?
解决SQLAlchemy异步分块查询的异步生成器实现问题
你遇到的TypeError: object async_generator can't be used in 'await' expression错误,核心原因是误用await调用异步生成器——异步生成器需要通过async for遍历,而非直接用await获取结果。以下是适配Python 3.11、SQLAlchemy 2.0.13的正确实现方案:
正确的BaseDAO泛型类实现
from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from typing import Generic, TypeVar, AsyncGenerator # 定义泛型类型变量 ModelType = TypeVar("ModelType") class BaseDAO(Generic[ModelType]): def __init__(self, model: type[ModelType]): self.model = model async def get_all_by_chunk( self, session: AsyncSession, chunk_size: int = 100 ) -> AsyncGenerator[list[ModelType], None]: offset = 0 while True: # 执行异步分块查询 result = await session.execute( select(self.model).offset(offset).limit(chunk_size) ) # 提取查询结果的标量列表 chunk = result.scalars().all() # 没有数据时终止生成器 if not chunk: break # 产出当前数据块 yield chunk offset += chunk_size
正确调用方式
调用异步生成器时,必须使用async for循环遍历产出的块,不能直接用await:
# 示例:替换YourModel为你的实际数据库模型 async def process_data(): async with AsyncSession(engine) as session: dao = BaseDAO(YourModel) # 用async for遍历生成器的每个块 async for chunk in dao.get_all_by_chunk(session, chunk_size=200): # 处理单条数据 for item in chunk: print(item.id)
关键要点说明
- 异步生成器通过
async def定义,返回类型标注为AsyncGenerator[生成的元素类型, 终止时返回值],这里我们每次产出一个模型列表,终止时返回None - 分块逻辑依赖SQL的
OFFSET和LIMIT,每次查询后偏移量累加,直到查询结果为空时退出循环 - 必须用
async for遍历生成器,这是Python异步生成器的标准调用方式,直接用await会触发类型错误
内容的提问来源于stack exchange,提问作者Seliverstov
相关产品推荐
相关产品推荐

