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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 07:19:53