FastAPI场景下如何在独立asyncio事件循环中使用Beanie?
我要开发一个集成Kafka与MongoDB的FastAPI应用,流程为:提供API将数据存入MongoDB并向Kafka发布事件,随后监听Kafka事件更新已有记录。我使用AIOKafka库,为避免Kafka监听器占用主应用事件循环,决定新开线程运行独立asyncio循环处理更新,但出现「任务绑定到不同事件循环」的错误。
错误信息
error occurred. Details: Task <Task pending name='Task-8' coro=<run_async_task..wrapper() running at fastApiMicroServices/proj/app/core/background.py:15> cb=[run_until_complete..done_cb()]> got Future <Future pending cb=[_chain_future.._call_check_cancel() at /Library/Frameworks/Python.framework/Versions/3.11/lib/python3.11/asyncio/futures.py:387]> attached to a different loop
相关代码片段
lifespan 函数
async def lifespan(application: FastAPI): logger = setup_root_logger() client = await init() session = await client.start_session() start_background_task(create_and_run_consumer, session=session, logger=logger) async with create_producer() as kafka_producer: application.producer = kafka_producer yield
init 函数
async def init() -> AsyncIOMotorClient: client = AsyncIOMotorClient(str(settings.MONGO_DB_URI), serverSelectionTimeoutMS=2000) await init_beanie(database=getattr(client, settings.MONGO_DB_NAME), document_models=gather_documents()) return client
start_background_task 函数
def start_background_task(coroutine, **kwargs): thread = threading.Thread( target=run_async_task, args=(coroutine,), kwargs=kwargs ) thread.start() return thread
run_async_task 函数
def run_async_task(coroutine, **kwargs): logger = kwargs.get('logger') loop = None async def wrapper(): try: if 'session' not in kwargs: kwargs['session'] = get_session() await coroutine(**kwargs) except Exception as e: logger.error(f"Error in background task: {e}") try: # loop = asyncio.new_event_loop() # asyncio.set_event_loop(loop) # loop.run_until_complete(wrapper()) asyncio.run(wrapper()) finally: if loop: loop.close()
create_and_run_consumer 函数
async def create_and_run_consumer(**kwargs): consumer = None logger = kwargs.get('logger') session = kwargs.get('session') try: consumer = AIOKafkaConsumer( settings.CONSUMER_TOPIC, bootstrap_servers=settings.KAFKA_DSN.split(','), group_id=settings.CONSUMER_GROUP) await consumer.start() logger.info('CONSUMER THREAD: Kafka consumer has started.') async for msg in consumer: logger.info(f'CONSUMER THREAD: consumed event for [{msg.key}] with payload: [{msg.value}]') payload = json.loads(msg.value) await Entity.find_one(Entity == msg.partition, session=session).update(Set({Entity.objects: payload['objects']})) except Exception as e: logger.error(f'CONSUMER THREAD: error occurred. Details: {str(e)}') finally: logger.info('CONSUMER THREAD: Kafka consumer has been closed.') if consumer: await consumer.stop()
问题根源
错误核心是:从主事件循环创建的MongoDB会话(session)被传到子线程的独立事件循环中使用。Beanie和Motor的会话、连接池对象都绑定了创建时的事件循环,跨循环调用必然触发「任务绑定到不同事件循环」的错误。
修复步骤
1. 子线程独立初始化MongoDB与Beanie
不要复用主线程的MongoDB资源,让子线程自己初始化连接、Beanie并创建专属会话:
def run_async_task(coroutine, **kwargs): logger = kwargs.get('logger') async def wrapper(): client = None session = None try: # 子线程独立初始化MongoDB连接和Beanie client = AsyncIOMotorClient(str(settings.MONGO_DB_URI), serverSelectionTimeoutMS=2000) await init_beanie(database=getattr(client, settings.MONGO_DB_NAME), document_models=gather_documents()) # 创建子线程循环专属的会话 session = await client.start_session() kwargs['session'] = session kwargs['client'] = client await coroutine(**kwargs) except Exception as e: logger.error(f"Error in background task: {e}") finally: # 清理子线程的MongoDB资源 if session: await session.end_session() if client: client.close() asyncio.run(wrapper())
2. 主线程不再传递session参数
修改lifespan,仅初始化主线程自己的MongoDB连接(供API使用),启动后台任务时不再传递主线程的session:
async def lifespan(application: FastAPI): logger = setup_root_logger() # 主线程初始化自身的MongoDB连接,供API接口使用 await init() # 启动后台任务,仅传递logger,不再传session start_background_task(create_and_run_consumer, logger=logger) async with create_producer() as kafka_producer: application.producer = kafka_producer yield
3. 消费者函数使用子线程专属会话
确保消费者逻辑使用子线程自己创建的会话执行数据库操作:
async def create_and_run_consumer(**kwargs): consumer = None logger = kwargs.get('logger') session = kwargs.get('session') # 现在session属于子线程的事件循环 try: consumer = AIOKafkaConsumer( settings.CONSUMER_TOPIC, bootstrap_servers=settings.KAFKA_DSN.split(','), group_id=settings.CONSUMER_GROUP) await consumer.start() logger.info('CONSUMER THREAD: Kafka consumer has started.') async for msg in consumer: logger.info(f'CONSUMER THREAD: consumed event for [{msg.key}] with payload: [{msg.value}]') payload = json.loads(msg.value) # 使用子线程专属会话执行更新 await Entity.find_one(Entity == msg.partition, session=session).update(Set({Entity.objects: payload['objects']})) except Exception as e: logger.error(f'CONSUMER THREAD: error occurred. Details: {str(e)}') finally: logger.info('CONSUMER THREAD: Kafka consumer has been closed.') if consumer: await consumer.stop()
关键注意事项
- 异步资源不能跨事件循环共享:Motor/Beanie的连接、会话、Future对象都绑定到创建它们的事件循环,跨线程循环使用必然报错。
- 子线程需独立初始化异步依赖:每个线程的asyncio循环都要有自己的MongoDB连接池和Beanie实例,不能复用主线程资源。
asyncio.run()自动管理循环:无需手动创建/设置事件循环,asyncio.run()会自动为子线程创建独立循环并在任务结束后销毁。
内容的提问来源于stack exchange,提问作者Vladyslav

