FastAPI爬虫API中await queue.get()阻塞致异步消费者挂起
问题描述
使用异步队列时遇到阻塞问题:调用await queue.get()会导致应用其余部分挂起。我正在开发基于FastAPI的网页抓取API,允许用户提交URL并批量处理,由消费者任务中的网页抓取器执行。但消费者等待队列新批次时,应用会停止运行——具体来说,我在FastAPI生命周期事件中调用start_consumers()启动应用,但应用因阻塞从未完成初始化。
已经了解asyncio.Queue.get()本身就是阻塞式的异步调用,但不清楚如何适配FastAPI场景,让消费者等待队列时不阻塞应用初始化。
原代码
import asyncio from asyncio.queues import Queue from internal.services.webscraping.base import BaseScraper class QueueController: _logger = create_logger(__name__) def __init__( self, scrapers: list[BaseScraper], batch_size: int = 50 ): self.queue = Queue() self.batch_size = batch_size self.scrapers = scrapers self.consumers = [] self.running = False async def put(self, urls: list[str]) -> None: """ 批量添加URL到队列 """ # 将URL列表按batch_size拆分后加入队列,队列满时等待空间 for i in range(0, len(urls), self.batch_size): batch = urls[i:i + self.batch_size] await self.queue.put(batch) async def get(self) -> list[str]: """ 从队列获取一批URL """ # 队列空时等待新项,获取到后返回 return await self.queue.get() async def consumer(self, scraper: BaseScraper) -> None: """ 消费者协程:处理队列中的URL批次 """ # 持续运行直到self.running为False,循环从队列取批次并处理 while self.running: try: batch = await self.get() if batch: records = await scraper.run(batch) # TODO: 处理结果存储 except Exception as e: # TODO: 添加完善的错误处理 ... raise e async def start_consumers(self) -> None: """ 启动消费者任务 """ self.running = True self.consumers = [ asyncio.create_task(self.consumer(scraper)) for scraper in self.scrapers ] await asyncio.gather(*self.consumers) async def stop_consumers(self) -> None: """ 优雅停止所有消费者任务 """ self.running = False for task in self.consumers: task.cancel()
解决方案
问题根源在start_consumers()方法里的await asyncio.gather(*self.consumers)——asyncio.gather()会等待所有传入的任务执行完毕才会返回,但消费者任务是无限循环(while self.running),这会导致FastAPI的初始化流程被永远阻塞。
修改方案很简单:移除await asyncio.gather(*self.consumers),只创建任务并保存引用即可。因为asyncio.create_task()会将协程加入事件循环后台运行,不会阻塞当前协程。
同时,为了保证消费者能正确响应停止信号,还需要在consumer协程中处理任务取消的异常,避免报错。
修改后的代码
import asyncio from asyncio.queues import Queue from internal.services.webscraping.base import BaseScraper class QueueController: _logger = create_logger(__name__) def __init__( self, scrapers: list[BaseScraper], batch_size: int = 50 ): self.queue = Queue() self.batch_size = batch_size self.scrapers = scrapers self.consumers = [] self.running = False async def put(self, urls: list[str]) -> None: """ 批量添加URL到队列 """ for i in range(0, len(urls), self.batch_size): batch = urls[i:i + self.batch_size] await self.queue.put(batch) async def get(self) -> list[str]: """ 从队列获取一批URL """ return await self.queue.get() async def consumer(self, scraper: BaseScraper) -> None: """ 消费者协程:处理队列中的URL批次 """ try: while self.running: batch = await self.get() if batch: records = await scraper.run(batch) # TODO: 处理结果存储 # 标记任务完成(可选,若队列需要跟踪未完成任务) self.queue.task_done() except asyncio.CancelledError: # 捕获任务取消信号,优雅退出 self._logger.info(f"消费者任务已取消:{scraper.__class__.__name__}") except Exception as e: self._logger.error(f"消费者任务出错:{str(e)}", exc_info=True) # 根据需求决定是否重新抛出或处理 async def start_consumers(self) -> None: """ 启动消费者任务 """ self.running = True self.consumers = [ asyncio.create_task(self.consumer(scraper)) for scraper in self.scrapers ] # 移除await gather,让任务后台运行 async def stop_consumers(self) -> None: """ 优雅停止所有消费者任务 """ self.running = False for task in self.consumers: if not task.done(): task.cancel() # 等待所有任务处理完取消信号 await asyncio.gather(*self.consumers, return_exceptions=True)
关键修改点
- 移除
start_consumers中的await asyncio.gather:让消费者任务在后台异步运行,不阻塞FastAPI初始化。 - 添加
asyncio.CancelledError捕获:消费者任务被取消时能优雅退出,避免未处理异常。 - 可选:添加
self.queue.task_done():如果需要跟踪队列中未完成的任务(比如后续要等待所有任务处理完再关闭应用),这个调用很有必要。 stop_consumers中添加await asyncio.gather:等待所有消费者任务处理完取消信号,确保优雅停止。
内容的提问来源于stack exchange,提问作者jda5
相关产品推荐
相关产品推荐

