FastAPI服务中如何用APScheduler正确调度异步协程函数?
正确使用APScheduler调度异步协程函数的方案
核心问题原因
直接将异步方法send_email传入AsyncIOScheduler时,调度器默认会把它当成普通函数执行,导致协程对象被创建但从未被await,触发警告。
解决方案步骤
1. 让AsyncIOScheduler复用FastAPI的事件循环
FastAPI运行在专属的asyncio事件循环中,需让调度器绑定这个循环,避免创建独立循环导致协程无法被正确调度:
from fastapi import FastAPI from apscheduler.schedulers.asyncio import AsyncIOScheduler import asyncio app = FastAPI() scheduler = None @app.on_event("startup") async def startup_event(): global scheduler # 获取当前运行的FastAPI事件循环 loop = asyncio.get_running_loop() # 初始化调度器并绑定循环 scheduler = AsyncIOScheduler(event_loop=loop) # 实例化状态机类 state_machine = ExchangeStateMachine() # 直接添加异步任务,调度器会自动处理await逻辑 scheduler.add_job(state_machine.send_email, "interval", seconds=60) # 替换为你的调度规则 scheduler.start()
2. 移除不必要的asyncio.run_coroutine_threadsafe调用
AsyncIOScheduler本身就是为异步环境设计的,无需用线程安全方法提交协程,直接传入异步方法即可。
3. 确保异步方法定义规范
检查send_email方法的异步语法,确保包含完整的异步逻辑:
class ExchangeStateMachine: def __init__(self, db_pool): self.db_pool = db_pool # 确保数据库连接池是异步初始化的 async def send_email(self): # 执行异步数据库操作 async with self.db_pool.acquire() as conn: async with conn.cursor() as cur: await cur.execute("SELECT email_content FROM task_queue WHERE status = 'pending'") result = await cur.fetchone() # 处理结果并发送邮件(邮件发送也建议用异步库,如aiosmtplib)
4. 调度器的优雅关闭
添加FastAPI关闭事件,避免资源泄漏:
@app.on_event("shutdown") async def shutdown_event(): if scheduler: scheduler.shutdown()
关键注意事项
- 所有异步资源(数据库连接池、邮件客户端等)必须和调度器共享同一个事件循环,避免跨循环调用导致的异常。
- 不要在同步代码块中直接调用异步方法,所有调度逻辑必须放在FastAPI的异步启动事件中完成。
内容的提问来源于stack exchange,提问作者jda5
相关产品推荐
相关产品推荐

