欢迎光临
我们一直在努力

Python 异步Web爬虫框架6道编程题:6 道编程题从入门到精通

本文为《Python全栈修炼之路》第22篇《阶段实战——异步Web爬虫框架的架构设计与底层实现》的配套练习。 六道题目由易到难,覆盖 asyncio事件循环、优先级队列调度、布隆过滤器去重、责任链中间件、生成器驱动数据流、令牌桶限速等核心知识点。每道题均包含完整题目描述、详细解题思路、关联知识点和可直接运行的参考代码。


题目一:asyncio基础——并发下载多个URL

题目描述

使用 asyncio 和 aiohttp 实现一个异步URL下载器,要求:

  • 定义一个异步函数 fetch_url(session, url),使用 aiohttp 发送GET请求,返回响应文本(前200个字符)
  • 定义一个异步函数 download_all(urls, max_concurrent),并发下载所有URL
  • 使用 asyncio.Semaphore 限制最大并发数为 max_concurrent
  • 记录每个URL的下载耗时,最终返回 {url: (status, content_preview, elapsed_time)}
  • 处理超时(timeout=10秒)和异常(返回状态码0和错误信息)
  • 测试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 并发请求。

    核心设计:

  • aiohttp.ClientSession 管理连接池,复用TCP连接
  • asyncio.Semaphore 控制并发数,防止对目标服务器造成过大压力
  • asyncio.gather 并发执行多个协程,等待全部完成
  • asyncio.TimeoutError 处理超时,aiohttp.ClientError 处理网络异常
  • 关键理解:

    • 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调度器,要求:

  • 定义 Request 类,包含 url(str)、priority(int, 默认0,越小越优先)、meta(dict)
  • 定义 PriorityScheduler 类,使用 asyncio.PriorityQueue 存储请求
  • 实现 push(request) 方法,将请求加入队列
  • 实现 pop() 方法,返回优先级最高的请求(阻塞等待)
  • 实现 is_empty() 方法,判断队列是否为空
  • 实现 __len__() 方法,返回队列长度
  • 验证:按不同优先级插入10个请求,观察出队顺序
  • 额外要求:

    • 处理优先级相同的情况(使用计数器保证FIFO)
    • 支持 maxsize 限制队列容量,满时 push 返回 False

    解题思路

    本题考察 优先级队列 的实现和 asyncio 异步队列 的使用。

    核心设计:

  • asyncio.PriorityQueue 基于 heapq 实现,插入/弹出 O(log n)
  • 优先级元组 (priority, count, request):当优先级相同时,计数器保证FIFO
  • asyncio.PriorityQueue 的 put 是异步的(如果队列满会等待),但本题要求满时返回False,所以需要检查容量
  • get() 是阻塞的,如果队列为空会等待直到有元素
  • 关键理解:

    • 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<

    赞(0)
    未经允许不得转载:171主机测评 » Python 异步Web爬虫框架6道编程题:6 道编程题从入门到精通
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址