本文为《Python全栈修炼之路》第22篇《阶段实战——异步Web爬虫框架的架构设计与底层实现》的配套练习。 六道题目由易到难,覆盖 asyncio事件循环、优先级队列调度、布隆过滤器去重、责任链中间件、生成器驱动数据流、令牌桶限速等核心知识点。每道题均包含完整题目描述、详细解题思路、关联知识点和可直接运行的参考代码。
题目一:asyncio基础——并发下载多个URL
题目描述
使用 asyncio 和 aiohttp 实现一个异步URL下载器,要求:
测试URL列表:
urls = [
\”https://httpbin.org/get\”,
\”https://httpbin.org/delay/1\”,
\”https://httpbin.org/status/404\”,
\”https://httpbin.org/status/500\”,
\”https://httpbin.org/bytes/1024\”,
]
解题思路
本题考察 asyncio 基础 和 aiohttp 并发请求。
核心设计:
关键理解:
- await 挂起当前协程,让出事件循环,其他协程可以执行
- Semaphore 的 async with 确保获取和释放的原子性
- ClientSession 的 __aenter__/__aexit__ 自动管理连接池生命周期
关联知识点
| async/await | 定义和调用协程 |
| aiohttp.ClientSession | 异步HTTP会话,管理连接池 |
| asyncio.Semaphore | 信号量,控制并发数 |
| asyncio.gather | 并发执行多个协程 |
| asyncio.TimeoutError | 超时异常处理 |
参考代码
import asyncio
import aiohttp
import time
from typing import Dict, Tuple
async def fetch_url(
session: aiohttp.ClientSession,
url: str,
timeout: float = 10.0
) –> Tuple[int, str, float]:
\”\”\”
异步获取单个URL的内容
返回: (状态码, 内容预览, 耗时秒数)
\”\”\”
start = time.time()
try:
async with session.get(url, timeout=aiohttp.ClientTimeout(total=timeout)) as response:
text = await response.text()
preview = text[:200] + \”…\” if len(text) > 200 else text
elapsed = time.time() – start
return response.status, preview, elapsed
except asyncio.TimeoutError:
elapsed = time.time() – start
return 0, f\”[超时] 请求超过 {
timeout} 秒\”, elapsed
except aiohttp.ClientError as e:
elapsed = time.time() – start
return 0, f\”[网络错误] {
type(e).__name__}: {
e}\”, elapsed
except Exception as e:
elapsed = time.time() – start
return 0, f\”[未知错误] {
type(e).__name__}: {
e}\”, elapsed
async def download_all(
urls: list,
max_concurrent: int = 5
) –> Dict[str, Tuple[int, str, float]]:
\”\”\”
并发下载所有URL,限制最大并发数
返回: {url: (status, preview, elapsed)}
\”\”\”
semaphore = asyncio.Semaphore(max_concurrent)
results = {
}
async def fetch_with_limit(url):
\”\”\”带并发限制的下载\”\”\”
async with semaphore:
return url, await fetch_url(session, url)
# 创建共享session(连接池复用)
async with aiohttp.ClientSession() as session:
# 并发执行所有任务
tasks = [fetch_with_limit(url) for url in urls]
completed = await asyncio.gather(*tasks, return_exceptions=True)
for item in completed:
if isinstance(item, Exception):
# 处理 gather 中的异常
print(f\”任务异常: {
item}\”)
continue
url, (status, preview, elapsed) = item
results[url] = (status, preview, elapsed)
return results
# ========== 验证 ==========
async def main():
urls = [
\”https://httpbin.org/get\”,
\”https://httpbin.org/delay/1\”,
\”https://httpbin.org/status/404\”,
\”https://httpbin.org/status/500\”,
\”https://httpbin.org/bytes/1024\”,
]
print(\”=== 异步并发下载测试 ===\”)
start = time.time()
results = await download_all(urls, max_concurrent=3)
total_elapsed = time.time() – start
print(f\”\\n总耗时: {
total_elapsed:.2f} 秒\”)
print(f\”平均每个URL: {
total_elapsed / len(urls):.2f} 秒\”)
print(\”\\n结果:\”)
for url, (status, preview, elapsed) in results.items():
print(f\”\\n URL: {
url}\”)
print(f\” 状态: {
status}, 耗时: {
elapsed:.2f}s\”)
print(f\” 预览: {
preview[:100]}\”)
if __name__ == \”__main__\”:
asyncio.run(main())
题目二:优先级队列调度器实现
题目描述
实现一个支持优先级调度的URL调度器,要求:
额外要求:
- 处理优先级相同的情况(使用计数器保证FIFO)
- 支持 maxsize 限制队列容量,满时 push 返回 False
解题思路
本题考察 优先级队列 的实现和 asyncio 异步队列 的使用。
核心设计:
关键理解:
- heapq 是最小堆,所以 priority 越小越优先
- 计数器解决优先级相同时的排序稳定性问题
- asyncio.Queue 的 maxsize 控制容量,qsize() 返回当前大小
关联知识点
| asyncio.PriorityQueue | 异步优先级队列,基于heapq |
| heapq | 最小堆实现,插入/弹出 O(log n) |
| 排序稳定性 | 计数器保证相同优先级FIFO |
| asyncio.Queue | 异步队列,支持阻塞put/get |
参考代码
import asyncio
from dataclasses import dataclass, field
from typing import Optional, Dict, Any
@dataclass(order=True)
class Request:
\”\”\”请求对象(支持优先级排序)\”\”\”
# 使用 dataclass 的排序功能,priority 越小越优先
priority: int = 0
url: str = field(compare=False)
meta: Dict[str, Any] = field(default_factory=dict, compare=False)
def __post_init__(self):
# 确保 priority 参与排序
pass
class PriorityScheduler:
\”\”\”优先级URL调度器\”\”\”
def __init__(self, maxsize: int = 0):
\”\”\”
maxsize: 队列最大容量,0表示无限制
\”\”\”
self._queue = asyncio.PriorityQueue(maxsize=maxsize)
self._count = 0 # 计数器,保证相同优先级FIFO
self._maxsize = maxsize
async def push(self, request: Request) –> bool:
\”\”\”
将请求加入队列
返回: True成功,False队列已满
\”\”\”
# 检查容量(非阻塞)
if self._maxsize > 0 and self._queue.qsize() >= self._maxsize:
return False
# 使用计数器保证相同优先级时的FIFO顺序
# PriorityQueue比较元组 (priority, count, request)
item = (request.priority, self._count, request)
self._count += 1
await self._queue.put(item)
return True
async def pop(self) –> Optional[Request]:
\”\”\”
获取优先级最高的请求
阻塞等待直到有请求可用
\”\”\”
priority, count, request = await self._queue.get()
return request
def is_empty(self) –> bool:
\”\”\”判断队列是否为空\”\”\”
return self._queue.empty()
def __len__(self) –> int:
\”\”\”返回队列长度\”\”\”
return self._queue.qsize()
@property
def maxsize(self) –> int:
return self._maxsize
# ========== 验证 ==========
async def main():
print(\”=== 优先级调度器测试 ===\”)
scheduler = PriorityScheduler(maxsize=100)
# 插入不同优先级的请求
requests = [
Request(url=\”https://example.com/low1\”, priority=10),
Request(url=\”https://example.com/high1\”, priority=1),
Request(url=\”https://example.com/mid1\”, priority=5),
Request(url=\”https://example.com/high2\”, priority=1), # 同优先级,先插入的先出
Request(url=\”https://example.com/low2\”, priority=10),
Request(url=\”https://example.com/mid2\”, priority=5),
Request(url=\”https://example.com/critical\”, priority=0),
Request(url=\”https://example.com/normal\”, priority=7),
]
print(f\”\\n插入 {
len(requests)} 个请求:\”)
for req in requests:
success = await scheduler.push(req)
print(f\” [{
req.priority:2d}] {
req.url.split(\’/\’)[–1]} -> {
\’OK\’ if success else \’FAIL\’}\”)
print(f\”\\n队列长度: {
len(scheduler)}\”)
print(\”\\n出队顺序(应从小到大):\”)
order = 1
while not scheduler.is_empty():
req = await scheduler.pop()
print(f\” #{
order:2d}: [{
req.priority:2d}] {
req.url.split(\’/\’)[–1]}\”)
order += 1
# 测试容量限制
print(\”\\n=== 容量限制测试 ===\”)
small_scheduler = PriorityScheduler(maxsize=3)
for i in range(5):
req = Request(url=f\”https://example.com/{
i}\”, priority=i)
success = await small_scheduler.push(req)
print(f\” 插入 #{
i}: {
\’成功\’ if success else \’失败(队列满)\’}\”)
if __name__ == \”__main__\”:
asyncio.run(main<


