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())
关键修改点说明
异步队列的使用:
asyncio.Queue.get()会自动阻塞,直到队列中有新任务,不需要你手动轮询空队列+await asyncio.sleep(0),代码更简洁高效。- 每个任务处理完成后调用
task_done(),队列的join()方法会等待所有任务都被标记为完成,确保没有遗漏的任务。
哨兵值的正确传递:
- 当Worker收到哨兵
None时,先把哨兵放回队列,这样同类型的其他Worker也能收到哨兵并退出,避免部分Worker提前终止。 - 只有当所有生产任务(比如导航生成URL、请求页面生成HTML)完成后,才会向队列发送哨兵,确保所有真实任务都被处理。
- 当Worker收到哨兵
Worker生命周期管理:
- 用
asyncio.create_task()创建任务,替代asyncio.ensure_future()(后者在Python 3.7+已被前者替代)。 - 主流程中通过
await navigate_worker、await nav_urls.join()等步骤,严格控制任务完成顺序,确保所有爬虫流程走完后再终止Worker。
- 用
这样修改后,你就不用担心Worker因为速度不一致提前退出的问题了——所有任务都会被完整处理,直到最后一个URL的解析工作完成,Worker才会收到哨兵并退出。
内容的提问来源于stack exchange,提问作者user3541631
相关产品推荐
相关产品推荐

