基于FastAPI+Celery开发,Beanie能否执行同步查询?
Beanie同步查询方案(适配Celery场景)
Beanie本身是基于Motor的异步MongoDB ODM,没有原生同步查询API,但可以通过以下几种实用方案在Celery这类同步任务环境中操作数据库:
方案1:封装异步逻辑执行工具
你当前用asyncio.run包裹异步代码的思路是可行的,可进一步封装成复用函数减少重复代码:
import asyncio from beanie import Document class MyDocument(Document): name: str def run_beanie_sync(coroutine): """封装异步Beanie操作的同步执行函数""" return asyncio.run(coroutine) # Celery任务中调用 @app.task def update_document_task(old_name, new_name): # 同步执行查询 target_doc = run_beanie_sync( MyDocument.find_one(MyDocument.name == old_name) ) # 同步执行更新 if target_doc: run_beanie_sync( target_doc.update({"$set": {"name": new_name}}) ) return f"Updated: {old_name} -> {new_name}"
方案2:复用Beanie模型+原生pymongo同步驱动
Beanie的模型定义可以直接和pymongo同步驱动配合使用,完全规避异步逻辑,适合高频Celery任务:
from pymongo import MongoClient from beanie import Document class MyDocument(Document): name: str # 初始化同步MongoDB客户端(可复用全局连接) sync_client = MongoClient("mongodb://localhost:27017") db = sync_client[MyDocument.Meta.database_name] collection = db[MyDocument.Meta.collection_name] # Celery任务中使用同步操作 @app.task def get_document_count_task(): # 同步查询统计 count = collection.count_documents({}) # 如需转换为Beanie模型实例 sample_doc = collection.find_one() if sample_doc: beanie_doc = MyDocument(**sample_doc) return f"Total documents: {count}"
该方案无需创建额外事件循环,性能更稳定,适合长期运行的Celery任务。
方案3:配置Celery异步Worker(可选)
若你的Celery版本支持,可通过eventlet或gevent异步Worker直接运行Beanie异步代码,无需包裹:
- 安装依赖:
pip install eventlet - 启动Worker:
celery worker -A your_app -P eventlet - 任务中直接写异步代码:
@app.task async def async_celery_task(): doc = await MyDocument.find_one(MyDocument.name == "test") await doc.update({"$set": {"name": "updated"}}) return doc.name
注意:该方案需调整Worker并发模型,需测试兼容性。
内容的提问来源于stack exchange,提问作者Puneet Gupta
相关产品推荐
相关产品推荐

