ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Tornado 从了解到精通(三):异步生产者 - 消费者模式深度解析

2026/8/14 3:17:04 拓冰建站 浏览量
Tornado 从了解到精通(三):异步生产者 - 消费者模式深度解析 前言在异步编程体系中生产者 - 消费者模式是实现任务分发、并发限流、批量处理的核心架构方案。Tornado 原生提供了tornado.queues模块功能与asyncio.Queue高度相似对标 Python 标准库中为线程设计的queue模块专为协程场景做了适配是 Tornado 异步开发的高频工具。本篇我们将从核心原理到实战落地完整讲解 Tornado 异步队列的使用方法并通过一个经典的异步网络爬虫案例带大家吃透异步生产者 - 消费者模式的实现逻辑。一、异步队列的核心运行机制tornado.queues.Queue是异步队列的核心实现它完全适配协程的调度特性不会阻塞线程只会在条件不满足时挂起当前协程、让出控制权。1. 协程的挂起与唤醒逻辑消费者侧取数据调用Queue.get()时如果队列为空当前协程会主动挂起并让出 CPU 控制权进入等待状态直到其他协程向队列中放入新元素等待的协程才会被唤醒继续执行。生产者侧放数据如果队列设置了最大容量maxsize调用Queue.put()时若队列已满生产者协程会自动挂起等待直到队列被消费者取出元素、腾出空间生产者协程才会完成元素放入并继续执行。2. 未完成任务计数器Tornado 队列内置了一个未完成任务计数器是实现 “全量任务等待” 的关键设计规则如下计数器初始值为 0每执行一次put()操作计数器数值 1每调用一次task_done()方法代表一个任务处理完成计数器数值 -1调用join()方法的协程会持续阻塞直到计数器归零代表所有入队的任务都已处理完成二、实战案例异步并发网络爬虫我们以 Tornado 官方文档站的全站爬虫为例完整演示异步生产者 - 消费者模式的落地爬虫从基础 URL 出发通过固定数量的协程并发抓取页面、解析链接、将新链接入队直到所有同域名链接都处理完成。1. 整体实现思路初始化异步队列放入起始 URL 作为初始任务启动固定数量的工作协程控制并发数作为消费者循环从队列中取出 URL 执行抓取每个协程抓取页面后解析页面内的所有链接将符合域名规则的新链接作为生产者放入队列每处理完一个 URL调用task_done()减少未完成任务计数主协程通过join()等待所有任务处理完成最终统计结果并结束程序2. 完整实现代码#!/usr/bin/env python3 import time from html.parser import HTMLParser from urllib.parse import urljoin, urldefrag from tornado import httpclient, queues, ioloop # 爬取的基础域名 base_url http://www.tornadoweb.org/en/stable/ # 并发协程数量控制抓取并发度 concurrency 10 async def get_links_from_url(url): 抓取指定URL的页面解析并返回页面内的所有链接 自动去除URL锚点并转换为绝对路径 例如 gen.html#tornado.gen.coroutine 会转为 http://www.tornadoweb.org/en/stable/gen.html response await httpclient.AsyncHTTPClient().fetch(url) print(ffetched {url}) html response.body.decode(errorsignore) return [urljoin(url, remove_fragment(new_url)) for new_url in get_links(html)] def remove_fragment(url): 去除URL中的锚点部分避免重复抓取同一页面的不同锚点 pure_url, frag urldefrag(url) return pure_url def get_links(html): 通过HTMLParser解析HTML内容提取所有a标签的href链接 class URLSeeker(HTMLParser): def __init__(self): super().__init__() self.urls [] def handle_starttag(self, tag, attrs): href dict(attrs).get(href) if href and tag a: self.urls.append(href) url_seeker URLSeeker() url_seeker.feed(html) return url_seeker.urls async def main(): # 初始化无容量限制的异步队列 q queues.Queue() start time.time() # 分别记录抓取中、已抓取、抓取失败的URL用于去重和统计 fetching, fetched, dead set(), set(), set() async def fetch_url(current_url): 单个URL的抓取处理协程 # 去重已在抓取队列中的URL直接跳过 if current_url in fetching: return print(ffetching {current_url}) fetching.add(current_url) try: urls await get_links_from_url(current_url) fetched.add(current_url) except Exception as e: print(ffetch failed: {current_url}, error: {e}) dead.add(current_url) return finally: # 无论抓取成功失败都标记该任务已完成 q.task_done() # 将新发现的、属于base_url域名下的链接放入队列 for new_url in urls: if new_url.startswith(base_url) and new_url not in fetched and new_url not in fetching: await q.put(new_url) # 放入初始任务启动整个爬取流程 await q.put(base_url) # 启动指定数量的消费者协程持续从队列取任务执行 workers [] for _ in range(concurrency): async def worker(): while True: url await q.get() await fetch_url(url) workers.append(worker()) # 阻塞等待直到所有入队任务都标记为完成 await q.join() # 统计爬取结果 cost time.time() - start print(\n *30) print(f爬取任务全部完成) print(f成功抓取页面数: {len(fetched)}) print(f抓取失败页面数: {len(dead)}) print(f总耗时: {cost:.2f} 秒) print(*30) if __name__ __main__: ioloop.IOLoop.current().run_sync(main)3. 关键逻辑说明并发数控制通过concurrency变量固定 worker 协程数量避免无限制发起请求导致目标服务器封禁或本地资源耗尽这是异步爬虫的核心风控点。URL 去重机制通过fetching和fetched两个集合分别记录 “正在抓取” 和 “已抓取完成” 的 URL彻底避免重复抓取、循环入队的问题。任务闭环设计每个 URL 处理结束后无论成功或失败都会调用q.task_done()保证未完成计数器的准确性主协程通过q.join()等待全量任务完成实现优雅收尾。异常兼容处理对网络请求增加异常捕获将请求失败的 URL 计入dead集合避免单个页面报错导致整个爬虫程序崩溃。三、核心 API 速查表API 方法功能说明queues.Queue(maxsize0)创建异步队列maxsize为 0 代表无容量限制await queue.put(item)向队列放入元素队列已满时自动挂起协程等待await queue.get()从队列取出元素队列为空时自动挂起协程等待queue.task_done()标记一个任务处理完成未完成计数器 -1await queue.join()阻塞当前协程直到队列中所有任务都被标记完成写在最后异步生产者 - 消费者模式是 Tornado 异步开发中的经典架构除了网络爬虫之外还广泛应用在异步任务处理、接口并发限流、批量数据消费等多种业务场景中。掌握tornado.queues的核心机制能够帮助我们写出执行效率更高、鲁棒性更强的异步服务代码。后续我们会继续深入 Tornado 的异步生态讲解更多进阶开发技巧。本系列为 Tornado 从入门到精通教程持续更新中欢迎关注跟进后续内容。