项目:异步网页爬虫
本章把第 11 章的 asyncio 和第 14 章的并发模式合起来,做一个"看起来像真爬虫"的项目。
沙箱说明
为了安全,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())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 包一层:
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())- 超时
asyncio.wait_for(...):避免某个请求挂死拖垮全局 - 重试:网络抖动是常态,3 次重试 + 指数退避(
delay = base * 2**attempt)很常用 - 并发数限制
asyncio.Semaphore(N):别一次开 10000 个任务把对端打挂
8. 信号量限流
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())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 秒,统计"成功 / 失败 / 超时"各自多少。
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 写真正的工具与服务了。