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

异步生产者消费者队列get()持续阻塞问题修复咨询

问题:异步队列消费者卡在entry = await queue.get()无进展

我有格式为article_urls: list[str]的URL列表,尝试实现以下逻辑:

  • 多个生产者Worker访问这些URL,从网页中提取PDF URL放入下载队列
  • 多个消费者Worker从该下载队列取出URL,下载文件后上传至S3

但程序运行后,消费者始终卡在entry = await queue.get()处没有进展,相关代码及终端输出如下:

代码实现

def get_first_pass_or_none(inp, driver):
    for x in inp:
        try:
            return x(driver)
        except:
            pass
    return None


async def url_producer(download_queue, pages_queue, producer_id):
    first = True
    while True:
        try:
            article_url = await pages_queue.get()
            instance_driver = wd[producer_id]  # or any other webdriver
            instance_driver.get(article_url)
            article_id = str(uuid.uuid4())

            if first:
                await asyncio.sleep(2)
                instance_driver.find_element(
                    By.XPATH, '//*[@id="onetrust-close-btn-container"]/button').click()
                first = False

            article_title = get_first_pass_or_none([lambda instance_driver: instance_driver.find_element(
                By.XPATH, '//*[@id="documentTitle"]').text], instance_driver)
            author = get_first_pass_or_none([lambda instance_driver: re.sub(r"\([^()]*\)", "", instance_driver.find_element(
                By.XPATH, '//*[@id="authordiv"]/span[2]/span/a/strong')
                .text)], instance_driver)
            publication_info = get_first_pass_or_none([lambda instance_driver: instance_driver.find_element(
                By.XPATH, '//*[@id="authordiv"]/span[2]/span').text], instance_driver)
            publication_location = get_first_pass_or_none([lambda publication_info: re.findall(
                r'\[(.*?)\]', publication_info)[0]], instance_driver)
            publication_date = publication_info

            output_metadata[article_id] = {
                "title": article_title,
                "author": author,
                "location": publication_location,
                "date": publication_date
            }
            pdf_url = instance_driver.find_element(
                By.CLASS_NAME, 'pdf-download').get_attribute('href')
            await download_queue.put({
                "article_id": article_id,
                "pdf_url": pdf_url,
            })
            print(download_queue.qsize())
            pages_queue.task_done()
        except Exception as e:
            logger.debug(f"Error {e}")
            keyboard.wait(keys[producer_id])


async def pdf_downloader(queue, consumer_id):
    while True:
        try:
            print(f"pdf_downloader {consumer_id} waiting for queue")
            entry = await queue.get()
            print(f"pdf_downloader {consumer_id} got queue entry")
            article_id = entry['article_id']
            pdf_url = entry['pdf_url']
            response = requests.get(pdf_url)
            pdf_content = response.content
            object_key = f"{article_id}.pdf"
            s3.Bucket(bucket_name).put_object(Key=object_key, Body=pdf_content)
            queue.task_done()
        except Exception as e:
            logger.debug(f"Error {e}")


async def main():
    # Create a shared queue
    download_queue = asyncio.Queue()
    pages_queue = asyncio.Queue()
    for page_url in article_urls:
        pages_queue.put_nowait(page_url)

    # Create two producers and two consumers
    producers = [asyncio.create_task(
        url_producer(download_queue, pages_queue, i)) for i in range(1)]
    consumers = [asyncio.create_task(
        pdf_downloader(download_queue, i)) for i in range(8)]

    # Wait for the producers to finish
    await asyncio.gather(*producers)

    # Cancel the consumers
    for consumer in consumers:
        consumer.cancel()

终端输出

pdf_downloader 0 waiting for queue
pdf_downloader 1 waiting for queue
pdf_downloader 2 waiting for queue
pdf_downloader 3 waiting for queue
pdf_downloader 4 waiting for queue
pdf_downloader 5 waiting for queue
pdf_downloader 6 waiting for queue
pdf_downloader 7 waiting for queue
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16

问题原因及修复方案

1. 同步阻塞操作卡死事件循环

生产者中使用的同步Selenium WebDriver和keyboard.wait()都是阻塞IO操作,会直接占用asyncio事件循环,导致消费者任务完全无法调度执行。

修复:

将所有同步阻塞操作放到asyncio.to_thread()中,让其在单独线程运行,不占用事件循环:

# 替换instance_driver.get(article_url)
await asyncio.to_thread(instance_driver.get, article_url)

# 弹窗点击操作也用to_thread包裹
await asyncio.to_thread(
    instance_driver.find_element,
    By.XPATH, '//*[@id="onetrust-close-btn-container"]/button'
).click()

# 其他find_element、text获取等操作同理,都用to_thread包裹

同时删除keyboard.wait(keys[producer_id]),这个同步等待会彻底卡死事件循环,可换成日志记录后直接终止当前任务。

2. 同步HTTP请求阻塞消费者

消费者中使用的requests.get()是同步阻塞操作,会拖慢事件循环,导致队列处理停滞。

修复:

改用异步HTTP库aiohttp,同时将boto3的同步上传操作也用to_thread包裹:

import aiohttp

async def pdf_downloader(queue, consumer_id):
    async with aiohttp.ClientSession() as session:
        while True:
            print(f"pdf_downloader {consumer_id} waiting for queue")
            entry = await queue.get()
            print(f"pdf_downloader {consumer_id} got queue entry")
            article_id = entry['article_id']
            pdf_url = entry['pdf_url']
            
            async with session.get(pdf_url) as response:
                pdf_content = await response.read()
                object_key = f"{article_id}.pdf"
                # 同步上传操作放入线程
                await asyncio.to_thread(
                    s3.Bucket(bucket_name).put_object,
                    Key=object_key, Body=pdf_content
                )
            queue.task_done()

3. 生产者无限循环无法终止

生产者用while True循环,但没有判断队列是否为空,导致await asyncio.gather(*producers)会永远等待,消费者也没机会处理完队列任务。

修复:

修改循环条件,当页面队列为空时退出:

async def url_producer(download_queue, pages_queue, producer_id):
    first = True
    # 队列空时退出循环
    while not pages_queue.empty():
        try:
            article_url = await pages_queue.get()
            # ... 原有逻辑 ...
            pages_queue.task_done()
        except Exception as e:
            logger.debug(f"Error {e}")
            # 出错也要标记任务完成,避免队列卡住
            pages_queue.task_done()

4. 消费者取消时机错误

当前代码在生产者完成后直接取消消费者,但此时下载队列可能还有未处理的任务,应该先等待下载队列处理完毕再取消。

修复:

async def main():
    download_queue = asyncio.Queue()
    pages_queue = asyncio.Queue()
    for page_url in article_urls:
        pages_queue.put_nowait(page_url)

    producers = [asyncio.create_task(url_producer(download_queue, pages_queue, i)) for i in range(1)]
    consumers = [asyncio.create_task(pdf_downloader(download_queue, i)) for i in range(8)]

    # 等待生产者完成所有页面处理
    await asyncio.gather(*producers)
    # 等待下载队列所有任务处理完毕
    await download_queue.join()
    # 取消消费者并等待其结束
    for consumer in consumers:
        consumer.cancel()
    await asyncio.gather(*consumers, return_exceptions=True)

额外代码错误修复

publication_location的get_first_pass_or_none调用参数错误:lambda期望接收publication_info,但实际传入的是instance_driver,会导致报错,修复如下:

publication_location = get_first_pass_or_none(
    [lambda pub_info: re.findall(r'\[(.*?)\]', pub_info)[0]],
    publication_info  # 传入正确的参数
)

内容的提问来源于stack exchange,提问作者maxlastname

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 19:43:06