异步生产者消费者队列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

