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

Asyncio任务未正常结束:哨兵机制引发的网页爬虫异常

解决异步爬虫Worker提前终止的问题

看起来你遇到的核心问题是用deque作为异步任务队列时,因为worker执行速度不一致,导致哨兵值(None)被提前放入队列,让还没处理完任务的Worker提前退出。deque本身不是异步安全的,而且你手动轮询空队列的方式不仅效率低,还很难精准控制任务的完成顺序和Worker的生命周期。

其实你提到的asyncio.Queue才是异步场景下管理任务队列的正确选择——它内置了异步安全的get()、put()方法,还有task_done()和join()机制,能完美解决你现在的问题。我帮你重构一下代码,顺便解释关键修改点:

核心修改思路

  • 把所有deque替换为asyncio.Queue,去掉手动轮询空队列的逻辑(Queue.get()会自动阻塞等待新任务)
  • 用task_done()和join()跟踪任务完成状态,确保所有爬虫任务都处理完毕后,再向队列发送哨兵值终止Worker
  • 调整Worker的停止逻辑,确保每个Worker收到哨兵后,把哨兵放回队列传递给其他同类型Worker,再退出

重构后的完整代码

import asyncio
import httpx
from bs4 import BeautifulSoup

domain_url = "https://www.example.com"  # 替换为你的实际域名

async def gather_search_links(html_sources, detail_urls):
    while True:
        data = await html_sources.get()
        if data is None:
            # 收到哨兵,放回队列让其他同类型Worker也能收到
            await html_sources.put(None)
            break
        soup = BeautifulSoup(data, "html.parser")
        results = soup.find_all("div", {"data-component": "search-result"})
        for record in results:
            atag = record.h2.a
            url = f"{domain_url}{atag.get('href')}"
            await detail_urls.put(url)
            print(f"已添加详情URL,当前总数: {detail_urls.qsize()}")
        # 标记当前任务完成
        html_sources.task_done()

async def get_page_source(urls, html_sources):
    async with httpx.AsyncClient() as client:
        while True:
            url = await urls.get()
            if url is None:
                # 收到哨兵,放回队列让其他同类型Worker也能收到
                await urls.put(None)
                break
            print(f"正在请求页面: {url}")
            response = await client.get(url)
            await html_sources.put(response.text)
            await asyncio.sleep(8)
            # 标记当前任务完成
            urls.task_done()
        # 所有URL处理完后,给解析队列发哨兵
        await html_sources.put(None)

async def navigate(urls):
    # 生成分页URL
    for page in range(2, 7):
        url = f"{domain_url}/?page={page}"
        await urls.put(url)
        await asyncio.sleep(0)
    # 所有分页URL生成完毕,给请求队列发哨兵
    await urls.put(None)

async def main():
    # 初始化异步队列
    nav_urls = asyncio.Queue()
    nav_html = asyncio.Queue()
    products_url = asyncio.Queue()

    # 创建Worker任务
    navigate_worker = asyncio.create_task(navigate(nav_urls))
    page_source_workers = [asyncio.create_task(get_page_source(nav_urls, nav_html)) for _ in range(2)]
    product_urls_workers = [asyncio.create_task(gather_search_links(nav_html, products_url)) for _ in range(1)]

    # 等待导航任务完成(所有分页URL生成完毕)
    await navigate_worker
    # 等待所有URL请求任务完成
    await nav_urls.join()
    # 等待所有HTML解析任务完成
    await nav_html.join()

    # 等待所有Worker退出
    await asyncio.gather(*page_source_workers, *product_urls_workers)

    # 可以在这里处理最终的products_url队列
    print(f"最终获取到的详情URL总数: {products_url.qsize()}")

if __name__ == "__main__":
    asyncio.run(main())

关键修改点说明

  1. 异步队列的使用:

    • asyncio.Queue.get()会自动阻塞,直到队列中有新任务,不需要你手动轮询空队列+await asyncio.sleep(0),代码更简洁高效。
    • 每个任务处理完成后调用task_done(),队列的join()方法会等待所有任务都被标记为完成,确保没有遗漏的任务。
  2. 哨兵值的正确传递:

    • 当Worker收到哨兵None时,先把哨兵放回队列,这样同类型的其他Worker也能收到哨兵并退出,避免部分Worker提前终止。
    • 只有当所有生产任务(比如导航生成URL、请求页面生成HTML)完成后,才会向队列发送哨兵,确保所有真实任务都被处理。
  3. Worker生命周期管理:

    • 用asyncio.create_task()创建任务,替代asyncio.ensure_future()(后者在Python 3.7+已被前者替代)。
    • 主流程中通过await navigate_worker、await nav_urls.join()等步骤,严格控制任务完成顺序,确保所有爬虫流程走完后再终止Worker。

这样修改后,你就不用担心Worker因为速度不一致提前退出的问题了——所有任务都会被完整处理,直到最后一个URL的解析工作完成,Worker才会收到哨兵并退出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 16:18:11