FastAPI+SQLAlchemy异步架构中MissingGreenlet错误的排查与修复咨询
FastAPI+SQLAlchemy异步架构中MissingGreenlet错误的排查与修复咨询
我现在遇到了一个异步架构下的SQLAlchemy报错问题,想请大家帮忙排查下。先给大家展示我的代码结构和问题细节:
1. Patch更新接口路由
我写了一个用于更新的Patch接口,代码如下:
@router.patch("/{id}", response_model=ResponseBlockSchema, status_code=status.HTTP_200_OK) async def update_block_router( update_data: UpdateBlockSchema, block_id: str = Query(..., alias="id"), block_service: BlockService = Depends(get_block_service), ): block = await block_service.update_block( block_id=uuid.UUID(block_id), update_data=update_data.dict(), ) if not block: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail="Block not found", ) return block
2. Pydantic请求/响应模型
对应的请求和响应Pydantic模型如下:
class BaseBlockSchema(BaseModel): model_config = ConfigDict(from_attributes=True) type: Optional[str] props: dict order: int class ChildBlockSchema(BaseBlockSchema): pass class ParentBlockSchema(BaseBlockSchema): pageId: Optional[uuid.UUID] parentId: Optional[uuid.UUID] class CreateBlockSchema(ParentBlockSchema): children: list[ChildBlockSchema] class UpdateBlockSchema(ParentBlockSchema): pass class ResponseBlockSchema(ParentBlockSchema): id: uuid.UUID class ResponseBlockSchemaAfterCreate(ResponseBlockSchema): children: list[ChildBlockSchema]
不用太在意UpdateBlockSchema,它和Patch方法的设计理念不太匹配,这个我之后会修复😅
3. 依赖注入:获取BlockService
BlockService的依赖注入逻辑:
async def get_block_service(session: AsyncSession = Depends(get_async_session)) -> BlockService: return BlockService( BlockRepository(session=session), )
4. 数据库会话配置
我的数据库异步配置是这样的:
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker, AsyncSession DATABASE_URL = f"postgresql+asyncpg://{config.POSTGRES_USER}:{config.POSTGRES_PASSWORD}@{config.POSTGRES_HOST}:{config.POSTGRES_PORT}/{config.POSTGRES_DB}" engine = create_async_engine(DATABASE_URL) async_session_maker = async_sessionmaker(engine, expire_on_commit=False) # 生成数据库会话 async def get_async_session() -> AsyncGenerator[AsyncSession, None]: async with async_session_maker() as session: yield session
我觉得异步会话的获取方式是正确的。
5. BlockService中的更新方法
这是核心的更新逻辑,问题就出在这里:
class BlockService: def __init__(self, block_repo: BlockRepository): self.block_repo = block_repo async def update_block(self, block_id: uuid.UUID, update_data: dict) -> Optional[Block]: block = await self.block_repo.get(block_id) if not block: return None for key, value in update_data.items(): if hasattr(block, key): setattr(block, key, value) if block.parentId == block.id: raise ValueError("Block cannot be its own parent.") return await self.block_repo.save(block)
在调试的时候,我发现在这个for循环部分触发了错误:
for key, value in update_data.items(): if hasattr(block, key): setattr(block, key, value)
我在循环之后的条件判断处打了断点,结果还没走到那里就报错了。
6. Repository实现
首先是BlockRepository:
class BlockRepository(BaseDBRepository[models.Block]): def __init__(self, session: AsyncSession, *args, **kwargs): super().__init__(models.Block, session, *args, **kwargs) async def get_blocks_by_page(self, page_id: UUID, many: bool = True) -> Optional[Sequence[models.Block] | models.Block]: return await self.get_where(many=many, whereclause=models.Block.page_id == page_id)
然后是基础的BaseDBRepository,主要用来处理数据库的常规交互:
AbstractModel = TypeVar('AbstractModel') def rollback_wrapper(func): @wraps(func) async def inner(self, *args, **kwargs): try: return await func(self, *args, **kwargs) except SQLAlchemyError as e: await self.session.rollback() logger.error(f"Error during DB operation: {e}") raise e return inner class BaseDBRepository(Generic[AbstractModel]): type_model: type[AbstractModel] def __init__(self, type_model: type[AbstractModel], session: AsyncSession, *args, **kwargs): self.type_model = type_model self.session = session @rollback_wrapper async def get(self, ident: Any, options: list = None ) -> AbstractModel | None: """Get an ONE model from the database with PK. :param options: :param ident: Key which need to find entry in database :return: """ async with self.session.begin_nested(): options = options or [] statement = select(self.type_model).options(*options).where(self.type_model.id == ident) result = await self.session.execute(statement) return result.unique().scalar_one_or_none() @rollback_wrapper async def get_where(self, many: bool = False, whereclause=None, limit: int | None = None, offset: int | None = None, order_by=None) -> Sequence[AbstractModel] | AbstractModel | None: """Get an ONE model from the database with whereclause. :param offset: :param many: :param whereclause: Clause by which entry will be found :param limit: Number of elements per query :param order_by: Name of field for ordering :return: Model if only one model was found, else None. """ async with self.session.begin_nested(): statement = select(self.type_model) if whereclause is not None: statement = statement.where(whereclause) if limit is not None: statement = statement.limit(limit) if offset is not None: statement = statement.offset(offset) if order_by is not None: statement = statement.order_by(order_by) result = await self.session.execute(statement) return result.unique().scalars().all() if many else result.scalar_one_or_none() @rollback_wrapper async def delete(self, obj: AbstractModel): logger.debug(f"[DB] Saving: {obj}") async with self.session.begin_nested(): await self.session.delete(obj) await self.session.commit() @rollback_wrapper async def save( self, obj: AbstractModel | Sequence[AbstractModel], many: bool = False ) -> AbstractModel | Sequence[AbstractModel]: """ :param obj: object or objects to save :param many: flag for many saves :return: """ async with self.session.begin_nested(): logger.debug(f"[DB] Saving: {obj}") if many and isinstance(obj, Sequence): self.session.add_all(obj) else: self.session.add(obj) await self.session.commit() if many and isinstance(obj, Sequence): for obj_item in obj: await self.session.refresh(obj_item) return obj
7. 数据库模型Block
对应的数据库模型:
class Block(Base): __tablename__ = "blocks" id = Column(UUID, primary_key=True, default=uuid.uuid4) type = Column(String, nullable=False) props = Column(JSON, nullable=False) parentId = Column(UUID, ForeignKey("blocks.id", ondelete="SET NULL"), nullable=True) pageId = Column(UUID, ForeignKey("pages.id", ondelete="SET NULL"), nullable=True) order = Column(Integer, nullable=False) created_at = Column(DateTime, default=func.now(), nullable=False) updated_at = Column(DateTime, default=func.now(), onupdate=func.now(), nullable=False) parent = relationship("Block", back_populates="children", remote_side="Block.id", lazy="subquery") children = relationship("Block", back_populates="parent", cascade="all, delete-orphan", lazy="subquery") page = relationship("Page", back_populates="blocks", lazy="subquery")
8. 报错信息
我收到的错误是:
raise exc.MissingGreenlet( sqlalchemy.exc.MissingGreenlet: greenlet_spawn has not been called; can't call await_only() here. Was IO attempted in an unexpected place?
备注:内容来源于stack exchange,提问作者Данил
相关产品推荐
相关产品推荐

