Learn
Python/18-project-async-scraper

项目:异步网页爬虫

本章把第 11 章的 asyncio 和第 14 章的并发模式合起来,做一个"看起来像真爬虫"的项目。

沙箱说明

⚠️本沙箱禁用了网络(--network none)

为了安全,Playground 跑在 Docker 容器里且禁止任何外部 HTTP 请求。所以本章:

  • "抓取"步骤用模拟数据演示并发模式(重点是 asyncio.gather 怎么用)
  • 真实项目里换成 aiohttp / httpx 即可,代码模式一模一样
  • /tmp 可写,所以本地落盘没问题

这个限制反而让示例自包含——你不需要任何 API key、网络环境就能跑。

1. 项目结构

scraper/
├── main.py          # 入口
├── fetcher.py       # 模拟/真实的网络抓取
├── parser.py        # 解析 HTML/JSON
├── aggregator.py    # 聚合统计
└── README.md

2. fetcher.py —— 异步抓取层

真实项目里用 aiohttp:

# fetcher.py —— 真实版本(需要 pip install aiohttp)
import aiohttp
import asyncio
 
async def fetch(session, url):
    async with session.get(url, timeout=10) as resp:
        resp.raise_for_status()
        return await resp.text()
 
async def fetch_all(urls):
    async with aiohttp.ClientSession() as session:
        tasks = [fetch(session, url) for url in urls]
        return await asyncio.gather(*tasks, return_exceptions=True)
💡沙箱里怎么写?

把 aiohttp 替换成"模拟网络":每次 await asyncio.sleep(secs) 一段随机时间,然后从一份内置字典里"取"出数据。接口形状完全一样——上层调用者根本不知道下面是真发请求还是假数据。

3. parser.py —— 解析

JSON API 通常不需要"解析",直接 json.loads;HTML 才需要 BeautifulSoup。这里我们处理 JSON-like 字典:

# parser.py
from typing import Any, Dict
 
def parse_post(raw: Dict[str, Any]) -> Dict[str, Any]:
    """从原始数据里挑出我们感兴趣的字段"""
    return {
        "id":   raw["id"],
        "title": raw["title"],
        "body":  raw["body"],
        "len":   len(raw["body"]),
    }

4. aggregator.py —— 聚合

# aggregator.py
from collections import Counter
from typing import List, Dict, Any
 
def summarize(posts: List[Dict[str, Any]]) -> Dict[str, Any]:
    if not posts:
        return {"count": 0, "avg_len": 0, "longest": None}
    lengths = [p["len"] for p in posts]
    longest = max(posts, key=lambda p: p["len"])
    return {
        "count":   len(posts),
        "avg_len": sum(lengths) / len(lengths),
        "longest": {"id": longest["id"], "title": longest["title"]},
    }

5. main.py —— 编排

# main.py
import asyncio
import json
from pathlib import Path
from fetcher import fetch_all
from parser import parse_post
from aggregator import summarize
 
URLS = [f"https://jsonplaceholder.typicode.com/posts/{i}" for i in range(1, 11)]
 
async def run():
    print("开始抓取", len(URLS), "条数据")
    raws = await fetch_all(URLS)
    posts = [parse_post(r) for r in raws if isinstance(r, dict)]
    summary = summarize(posts)
    print("汇总:", summary)
 
    out = Path("/tmp/scraper_result.json")
    out.write_text(
        json.dumps({"summary": summary, "posts": posts},
                   ensure_ascii=False, indent=2),
        encoding="utf-8"
    )
    print("已写入:", out)
 
if __name__ == "__main__":
    asyncio.run(run())

6. 可运行演示(沙箱版本)

把所有文件合并成一个 Playground,用模拟数据替代真请求:

异步爬虫:模拟数据演示
import asyncio
import json
from pathlib import Path
from typing import Any, Dict, List
 
# ---- 模拟的"远端数据库",对应 jsonplaceholder 的 posts ----
FAKE_DB = {
    i: {
        "id": i,
        "title": f"Post #{i}: 学习 Python 异步",
        # 用 (i % 4) + 1 让长度确定,方便演示
        "body":  "在沙箱里我们用 asyncio.sleep 模拟网络延迟。" * ((i % 4) + 1),
        "userId": (i % 3) + 1,
    }
    for i in range(1, 11)
}
 
async def fetch(url: str) -> Dict[str, Any]:
    """沙箱版 fetch:解析 URL 拿 id, 模拟网络延迟, 返回数据"""
    # URL 形如 https://jsonplaceholder.typicode.com/posts/3
    post_id = int(url.rstrip("/").split("/")[-1])
    await asyncio.sleep(0.05)          # 模拟 50ms 网络
    return FAKE_DB[post_id]
 
async def fetch_all(urls: List[str]) -> List[Any]:
    """asyncio.gather 并发抓取全部"""
    tasks = [fetch(u) for u in urls]
    return await asyncio.gather(*tasks, return_exceptions=True)
 
def parse_post(raw: Dict[str, Any]) -> Dict[str, Any]:
    return {"id": raw["id"], "title": raw["title"],
            "body": raw["body"], "len": len(raw["body"])}
 
def summarize(posts: List[Dict[str, Any]]) -> Dict[str, Any]:
    if not posts:
        return {"count": 0, "avg_len": 0, "longest": None}
    longest = max(posts, key=lambda p: p["len"])
    return {
        "count":   len(posts),
        "avg_len": sum(p["len"] for p in posts) / len(posts),
        "longest": {"id": longest["id"], "title": longest["title"]},
    }
 
async def main():
    urls = [f"https://jsonplaceholder.typicode.com/posts/{i}" for i in range(1, 11)]
    print(f"开始抓取 {len(urls)} 个 URL ...")
 
    raws = await fetch_all(urls)
    print("抓取完成 (并发 50ms × 1 轮)")
 
    # 区分成功 / 失败
    posts  = [parse_post(r) for r in raws if isinstance(r, dict)]
    failed = [r for r in raws if isinstance(r, BaseException)]
    print(f"成功 {len(posts)} 条, 失败 {len(failed)} 条")
 
    summary = summarize(posts)
    print("")
    print("--- 汇总 ---")
    print(f"  count:   {summary['count']}")
    print(f"  avg_len: {summary['avg_len']:.1f}")
    print(f"  longest: #{summary['longest']['id']} {summary['longest']['title']}")
 
    out = Path("/tmp/scraper_result.json")
    out.write_text(
        json.dumps({"summary": summary, "posts": posts}, ensure_ascii=False, indent=2),
        encoding="utf-8"
    )
    print("")
    print(f"已写入: {out} ({out.stat().st_size} bytes)")
 
asyncio.run(main())
ℹ️为什么 avg_len = 72.5?

body 长度是 29 * ((i % 4) + 1),10 条分别是 58/87/116/29/58/87/116/29/58/87。max() 返回第一次出现最大值的那条,所以 longest = #3(116)。手算 (58+87+116+29+58+87+116+29+58+87)/10 = 72.5。

7. 加上超时与重试

真实爬虫必须考虑慢请求、连接失败、限流。给 fetch 包一层:

带超时 + 重试的 fetch
import asyncio, random
from typing import Any, Dict, List
 
FAIL_RATE = 0.2
SLOW_RATE = 0.2
 
async def raw_fetch(url: str) -> Dict[str, Any]:
    """低层:可能慢、可能失败"""
    delay = random.uniform(0.05, 0.5) if random.random() < SLOW_RATE else random.uniform(0.01, 0.05)
    await asyncio.sleep(delay)
    if random.random() < FAIL_RATE:
        raise RuntimeError(f"网络抖动: {url}")
    return {"id": int(url.rsplit("/", 1)[1]), "ok": True}
 
async def fetch_with_retry(url: str, *, max_retries=3, timeout=0.2) -> Dict[str, Any]:
    """带超时和重试的高层 fetch"""
    last = None
    for i in range(max_retries):
        try:
            return await asyncio.wait_for(raw_fetch(url), timeout=timeout)
        except (asyncio.TimeoutError, RuntimeError) as e:
            last = e
            print(f"  {url} 第 {i+1} 次失败: {type(e).__name__}: {e}")
    raise last
 
async def main():
    random.seed(7)
    urls = [f"http://api.example.com/posts/{i}" for i in range(8)]
    results = await asyncio.gather(
        *(fetch_with_retry(u) for u in urls),
        return_exceptions=True
    )
    ok = [r for r in results if isinstance(r, dict)]
    bad = [r for r in results if isinstance(r, BaseException)]
    print(f"成功 {len(ok)} / 失败 {len(bad)}")
 
asyncio.run(main())
💡real-world 三件套
  1. 超时 asyncio.wait_for(...):避免某个请求挂死拖垮全局
  2. 重试:网络抖动是常态,3 次重试 + 指数退避(delay = base * 2**attempt)很常用
  3. 并发数限制 asyncio.Semaphore(N):别一次开 10000 个任务把对端打挂

8. 信号量限流

Semaphore 限制并发数
import asyncio, time
 
# 同时最多 3 个任务
sem = asyncio.Semaphore(3)
 
async def worker(i):
    async with sem:        # 抢到信号量才执行
        await asyncio.sleep(0.1)
        return f"done-{i}"
 
async def main():
    start = time.perf_counter()
    # 10 个任务,但同时只有 3 个在跑
    results = await asyncio.gather(*(worker(i) for i in range(10)))
    print(f"完成 {len(results)} 个, 总耗时 {time.perf_counter()-start:.2f}s")
 
asyncio.run(main())
ℹ️为什么是 0.4s?

10 个任务每个 0.1s,但同时只跑 3 个 → 至少要 ⌈10/3⌉ = 4 轮 = 0.4s。如果不限流,10 个并发跑就是 0.1s。并发太快反而容易触发反爬或被对方封 IP。

🎯 练习

把"带超时的 fetch"和"Semaphore 限流"组合起来:写一个 fetch_many(urls, concurrency=5, timeout=0.2) 函数,要求并发上限 5,每个请求超时 0.2 秒,统计"成功 / 失败 / 超时"各自多少。

fetch_many:限流 + 超时
import asyncio, random
from typing import List, Dict, Any
 
async def raw_fetch(url: str) -> Dict[str, Any]:
    await asyncio.sleep(random.uniform(0.05, 0.3))
    if random.random() < 0.2:
        raise RuntimeError("网络错误")
    return {"id": int(url.rsplit("/", 1)[1])}
 
async def fetch_many(urls: List[str], concurrency: int = 5, timeout: float = 0.2):
    # 你的实现
    pass
 
async def main():
    random.seed(1)
    urls = [f"http://api.example.com/{i}" for i in range(20)]
    result = await fetch_many(urls)
    print(result)
 
asyncio.run(main())
🎯提示
  • sem = asyncio.Semaphore(concurrency)
  • async def one(url): async with sem: return await asyncio.wait_for(raw_fetch(url), timeout)
  • results = await asyncio.gather(*(one(u) for u in urls), return_exceptions=True)
  • 分类:成功 = isinstance(r, dict);超时 = isinstance(r, asyncio.TimeoutError);其他 = isinstance(r, BaseException)

小结

  • ✅ asyncio.gather + return_exceptions=True 让你"全都要"且不互相干扰
  • ✅ 沙箱网络受限用模拟数据演示接口形状,真实项目换 aiohttp 即可
  • ✅ asyncio.wait_for 给单个请求加超时
  • ✅ 重试(指数退避)能扛住大部分网络抖动
  • ✅ asyncio.Semaphore(N) 限流,避免并发过大触发反爬
  • ✅ 完整的爬虫项目分层:fetcher(IO)→ parser(清洗)→ aggregator(统计)→ main(编排)

🎉 Python 18 章完结!你已经系统掌握了 Python 的核心语法、OOP、并发、测试和项目实战,可以开始用 Python 写真正的工具与服务了。