Celery RPC调用阻塞FastAPI并发问题及替代方案咨询
问题分析
你遇到的核心问题是Celery的task.get()是同步阻塞调用,会占用FastAPI异步事件循环的线程,导致事件循环无法处理其他请求,最终所有请求串行执行,响应时间远超预期。
解决方案
方案1:用asyncio.to_thread包装阻塞调用(Python 3.9+)
将task.get()放到单独的线程池中执行,不阻塞FastAPI的事件循环,让其他请求可以并行处理。
修改main.py中的search接口:
import asyncio from celery.exceptions import TimeoutError RESPONSE_TIMEOUT_SECONDS = 1.5 @app.get('/search/') async def search(a,b): # 提交任务到Celery,不阻塞 task = remote_procedure.apply_async( (a,b), expires=RESPONSE_TIMEOUT_SECONDS+0.1, retry=False ) try: # 把阻塞的get()放到线程池执行,事件循环可继续处理其他请求 res = await asyncio.to_thread( task.get, timeout=RESPONSE_TIMEOUT_SECONDS ) except TimeoutError: print('timeout') return {"res":"timeout"} return {"res":res}
方案2:用sync_to_async兼容低版本Python
如果你的Python版本低于3.9,使用asgiref.sync.sync_to_async实现同样效果:
from asgiref.sync import sync_to_async from celery.exceptions import TimeoutError RESPONSE_TIMEOUT_SECONDS = 1.5 @app.get('/search/') async def search(a,b): task = remote_procedure.apply_async( (a,b), expires=RESPONSE_TIMEOUT_SECONDS+0.1, retry=False ) try: res = await sync_to_async(task.get)( timeout=RESPONSE_TIMEOUT_SECONDS ) except TimeoutError: print('timeout') return {"res":"timeout"} return {"res":res}
额外优化建议
- 调整Celery Worker数量:远程GPU机器的Worker数量要匹配GPU并发能力,比如单GPU建议设置2-4个Worker,避免GPU过载导致任务排队。
- 检查Broker性能:确保RabbitMQ配置合理(如增加队列长度、调整连接数),避免消息传递延迟。
- 监控任务队列:用Celery Flower监控任务状态,排查是否有任务堆积或异常。
内容的提问来源于stack exchange,提问作者Nick Zorander
相关产品推荐
相关产品推荐

