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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 17:30:11