尧图精选

Python后端爬虫专题14:并发不是越大越快——asyncio任务、Semaphore与连接预算

🕒 发布时间:2026/9/28 8:50:27 📁 来源:尧图网络
Python后端爬虫专题14并发不是越大越快——asyncio任务、Semaphore与连接预算上一篇练习完整答案首次请求不带验证器服务端返回200 ETag: v1 body客户端保存正文与 v1。第二次发送If-None-Match: v1服务端确认未变化后返回304 ETag: v1没有 body。内容更新后同样请求 v1服务端返回200 ETag: v2 新 body客户端保存并更新业务字段。304 时空正文不是新版本生成空快照会丢失可重放证据计算指纹则会把“未传输”误判为“职位被清空”。1000 个详情中 950 个未变化每个 50 KB少传输950 × 50 47,500 KB按十进制约 47.5 MB实际节省还受压缩、响应头和协议影响。把并发想成一家小餐馆顺序下载像只有一名服务员第一桌等网络时其他桌没人照顾。无限asyncio.gather又像一次接待五百桌厨房、连接池和目标网站都会被压垮。合理做法是先确定容量某个来源同时最多 4 个详情请求数据库写入仍按稳定顺序处理HTTP 连接池的上限不低于任务并发但也不能无限。asyncio适合这里是因为下载大部分时间在等待网络而不是执行 CPU 计算。事件循环可在一个请求等待时推进另一个请求它不会让纯 Python 的重 CPU 工作自动并行。HTML 特别大或清洗算法很重时应测量后再考虑进程池或独立 Worker。bounded_map 的合同bounded_map(items, worker, limit4)接收一个输入序列和异步 worker返回与输入同顺序的结果列表。内部为每个输入建立协程但协程进入实际 worker 前必须拿到 Semaphoreasync with semaphore保证异常和取消时也释放许可。为什么保留顺序下载完成顺序受网络抖动影响如果报告与测试每次都不同调试困难。课程把网络阶段并发、数据库阶段顺序处理先得到(url, response, error)列表再逐项快照和 upsert。这样不在线程之间共享 SQLAlchemy Session也让错误报告按发现顺序稳定。cd project.\.venv\Scripts\python.exe-m pytest tests\test_concurrency.py::test_bounded_map_preserves_order_without_exceeding_concurrency-q检查点的 worker 会记录活动数和峰值并故意让不同项等待不同时间。最终既断言peak limit也断言输出顺序仍对应输入。只测总耗时容易在机器负载变化时误报直接测并发峰值更贴近合同。三个数字不要互相打架应用层detail_concurrency4限制单次 CrawlerHTTPXLimits(max_connections...)限制一个客户端的总连接Celery concurrency 决定一台机器同时运行多少采集任务。若 Worker 同时跑 8 个任务每个任务对同一域并发 4请求峰值可能是 32而不是 4。生产预算应从来源允许频率倒推。例如目标允许每秒 5 次且单请求 P95 约 400 ms同域并发 2 通常已接近吞吐上限再加大只会产生 429。多租户共享同一来源时要做跨任务、按域名的全局限速单个 Semaphore 不够。本项目用 Scrapy AutoThrottle 讲解自适应节奏但仍不把 429 当作“换 IP 再冲”的信号。异常、取消与超时各管一层worker 捕获每个详情的下载异常并把异常对象作为结果带回因此一个坏详情不会让gather取消全部任务。HttpFetcher 自己负责单请求超时与有限重试Celery 负责整个业务任务因基础设施失败而重跑。把所有异常都吞掉会让任务假成功把所有异常都抛出又会失去部分成功这就是CrawlReport.errors存在的原因。如果应用收到关闭信号正在等待 Semaphore 的协程会被取消。async with可释放已拿到的许可但数据库已提交的条目不会自动撤销幂等 upsert 让重启后重跑安全。取消语义最终仍要与任务队列的 late acknowledgement 配合第 21 篇会落到 Worker 故障恢复。如何调参而不是猜数字从 1 开始记录详情下载 P50/P95、429 比例、失败率、CPU、连接等待和目标服务公开限制逐级到 2、4吞吐不再明显提高或错误开始增加就回退。不要拿自己的千兆网络当服务端容量也不要把本地 TargetLab 的极低延迟直接搬到公网配置。本篇完整并发模块模块很短因为它只负责一个稳定职责限制进入 worker 的数量并保持结果次序。下载重试、数据校验、数据库事务都不塞进这个工具函数。有上限的异步并发工具。importasynciofromcollections.abcimportAwaitable,Callable,SequencefromtypingimportTypeVar InputTTypeVar(InputT)OutputTTypeVar(OutputT)asyncdefbounded_map(values:Sequence[InputT],worker:Callable[[InputT],Awaitable[OutputT]],*,limit:int,)-list[OutputT]:限制同时工作的协程数量并按输入顺序返回结果。iflimit1:raiseValueError(concurrency limit must be positive)semaphoreasyncio.Semaphore(limit)asyncdefguarded(value:InputT)-OutputT:asyncwithsemaphore:returnawaitworker(value)returnlist(awaitasyncio.gather(*(guarded(value)forvalueinvalues)))本篇课后练习编写一个完整异步脚本输入 0—9限制并发为 3worker 等待后返回平方并断言顺序与峰值。Worker concurrency6每任务详情并发4三个任务访问 A 域、三个任务访问 B 域算出总峰值和单域理论峰值。解释为什么不应在并发下载协程中共享同一个 SQLAlchemy Session 写库。下一篇会处理必须登录、但已明确获得授权的内部来源。
上一篇/下一篇内容由系统自动关联 返回资讯列表 →