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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 09:05:58