尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Python爬虫并发技术对比:多线程、多进程与协程性能分析

Python爬虫并发技术对比:多线程、多进程与协程性能分析 1. 爬虫并发技术选型背景当我们需要从互联网上批量获取数据时单线程的爬虫程序就像一个人在超市收银台排队结账 - 效率低得令人发指。这时候就需要引入并发技术让我们的爬虫能够分身处理多个请求。目前主流的三种并发方案各有特点多线程爬虫相当于雇佣了一群收银员线程但共用一个收银台GIL锁。虽然可以同时处理多个顾客请求但每次只能有一个收银员真正操作收银机。这种模式在I/O等待网络请求占主要时间的场景下效果不错因为收银员可以在等待顾客掏钱包时去服务下一位。多进程爬虫则是直接开了多个收银通道进程每个通道都有自己独立的收银系统。这种方式能真正并行工作但开设新通道的成本较高进程创建和切换开销大而且通道间的沟通进程间通信比较麻烦。异步协程爬虫更像是一个超级收银员他可以在处理一个顾客的间隙快速切换到另一个顾客而且切换几乎不花时间。这种模式特别适合顾客很多但每个交易都很简单的情况纯I/O密集型任务。2. 测试环境与方法论2.1 硬件与软件配置为了保证测试结果的可靠性我们搭建了统一的测试环境硬件平台Intel i7-12700H处理器14核20线程16GB DDR4内存1Gbps有线网络连接操作系统Ubuntu 22.04 LTSPython环境Python 3.10.12使用venv创建隔离环境关键库版本requests2.31.0aiohttp3.9.1urllib32.0.72.2 测试目标与指标我们选择httpbin.org的/get接口作为测试目标这个公开API可以稳定返回请求信息非常适合作为基准测试对象。每个测试案例执行1000次请求记录以下指标总耗时从第一个请求开始到最后一个请求完成的时间CPU利用率通过psutil库监控进程资源使用情况内存占用测试过程中的峰值内存使用量网络延迟平均每个请求的网络往返时间2.3 控制变量设计为确保测试公平性我们设置了以下控制条件禁用DNS缓存设置requests和aiohttp的DNS缓存TTL为0统一设置5秒请求超时关闭SSL验证避免证书验证带来的性能差异每个测试案例运行3次取平均值测试前预热网络连接先发送10个请求不计入统计3. 各方案实现细节3.1 同步基准实现作为性能对比的基线我们首先实现一个最简单的同步爬虫import requests import time def sync_fetch(url): try: response requests.get(url, timeout5) return response.status_code except Exception as e: return str(e) def sync_crawler(url, count): start time.perf_counter() results [sync_fetch(url) for _ in range(count)] elapsed time.perf_counter() - start print(f同步爬虫完成{count}次请求耗时{elapsed:.2f}秒) return elapsed, results这个实现简单直接但效率也是最低的 - 它就像固执的老人一定要等上一个请求完全结束才会开始下一个。3.2 多线程实现优化我们使用线程池优化多线程实现避免频繁创建销毁线程的开销from concurrent.futures import ThreadPoolExecutor import threading from collections import defaultdict class ThreadedCrawler: def __init__(self, max_workers10): self.executor ThreadPoolExecutor(max_workersmax_workers) self.counter defaultdict(int) self.lock threading.Lock() def fetch(self, url): try: response requests.get(url, timeout5) with self.lock: self.counter[response.status_code] 1 return response.status_code except Exception as e: with self.lock: self.counter[str(e)] 1 return str(e) def run(self, url, count): start time.perf_counter() futures [self.executor.submit(self.fetch, url) for _ in range(count)] results [f.result() for f in futures] elapsed time.perf_counter() - start print(f线程池({self.executor._max_workers} workers)完成{count}次请求耗时{elapsed:.2f}秒) print(状态码统计:, dict(self.counter)) return elapsed, results关键优化点使用ThreadPoolExecutor管理线程生命周期采用线程安全的defaultdict统计响应状态通过with语句确保锁的正确获取和释放3.3 多进程实现进阶多进程实现需要考虑进程间通信问题我们使用共享内存减少序列化开销from multiprocessing import Pool, Manager import os def process_fetch(args): url, result_dict args try: response requests.get(url, timeout5) result_dict[os.getpid()].append(response.status_code) return response.status_code except Exception as e: result_dict[os.getpid()].append(str(e)) return str(e) def process_crawler(url, count, processes8): manager Manager() result_dict manager.dict() for i in range(processes): result_dict[i] manager.list() start time.perf_counter() with Pool(processes) as pool: pool.map(process_fetch, [(url, result_dict)]*count) elapsed time.perf_counter() - start print(f进程池({processes} workers)完成{count}次请求耗时{elapsed:.2f}秒) return elapsed, result_dict特别注意事项Manager创建的共享数据结构有性能损耗每个进程维护自己的结果列表避免频繁锁竞争使用进程ID作为字典key方便结果归集3.4 异步协程深度优化异步协程实现我们增加了连接池和重试机制import aiohttp import asyncio from aiohttp import TCPConnector from tenacity import retry, stop_after_attempt, wait_exponential class AsyncCrawler: def __init__(self, conn_limit100): self.connector TCPConnector(limitconn_limit, force_closeTrue) self.session None async def __aenter__(self): self.session aiohttp.ClientSession(connectorself.connector) return self async def __aexit__(self, exc_type, exc, tb): await self.session.close() retry(stopstop_after_attempt(3), waitwait_exponential(multiplier1, min1, max5)) async def fetch(self, url): try: async with self.session.get(url, timeout5) as resp: return resp.status except Exception as e: print(fRequest failed: {str(e)}) raise async def crawl(self, url, count): start time.perf_counter() tasks [self.fetch(url) for _ in range(count)] results await asyncio.gather(*tasks, return_exceptionsTrue) elapsed time.perf_counter() - start print(f异步爬虫完成{count}次请求耗时{elapsed:.2f}秒) return elapsed, results高级特性使用TCPConnector控制最大并发连接数实现上下文管理协议确保Session正确关闭集成tenacity实现指数退避重试支持异常结果收集不中断整体任务4. 性能对比与分析4.1 原始数据对比我们在相同环境下运行四种实现得到如下数据爬虫类型总耗时(秒)CPU利用率内存峰值(MB)网络吞吐(MB/s)同步基准312.458%450.32多线程(10线程)35.6728%722.81多进程(8进程)29.8378%3353.36异步协程16.0435%556.244.2 关键发现解读异步协程优势明显在纯I/O场景下协程的轻量级切换使其性能达到多线程的2倍以上。这得益于用户态任务调度避免内核切换开销事件循环高效管理大量空闲连接零拷贝技术减少数据传输损耗多进程的代价虽然多进程的CPU利用率最高但内存占用是其他方案的4-6倍。这是因为每个Python进程都有独立的内存空间进程间通信需要数据序列化测试中使用的Manager进一步增加了开销线程池的平衡点在8核CPU上10个线程的配置表现出不错的性价比比进程方案节省75%内存性能损失控制在20%以内编程模型最简单直观4.3 各方案瓶颈分析通过perf工具分析我们发现不同实现的主要瓶颈点同步实现99%时间花费在socket.recv等待网络响应多线程实现约30%时间消耗在线程切换和GIL竞争多进程实现15-20%开销用于进程间通信和数据序列化异步实现主要瓶颈转移到DNS查询占时约8%5. 实战建议与陷阱规避5.1 选型决策树根据项目需求选择合适方案是否纯I/O密集型 ├─ 是 → 是否需要极高并发(10k QPS)? │ ├─ 是 → 选择异步协程 │ └─ 否 → 多线程更简单 └─ 否 → 计算密集型占比? ├─ 30% → 多线程线程池 ├─ 30-70% → 多进程进程池 └─ 70% → 考虑分布式方案5.2 性能调优技巧连接复用所有方案都应启用keep-alive# requests会话 session requests.Session() # aiohttp连接池 connector TCPConnector(keepalive_timeout60)智能限流根据响应时间动态调整并发度class AdaptiveLimiter: def __init__(self, max_concurrent100): self.semaphore asyncio.Semaphore(max_concurrent) self.last_response_time 0 async def throttle(self): if time.time() - self.last_response_time 1.0: self.semaphore._value min( self.semaphore._value 5, self.semaphore._initial_value ) else: self.semaphore._value max( self.semaphore._value - 2, 1 ) await self.semaphore.acquire()DNS缓存特别是异步方案中缓存DNS查询from aiodnsresolver import Resolver, cache resolver Resolver(cachecache.Cache()) connector TCPConnector(resolverresolver)5.3 常见陷阱与解决方案TCP端口耗尽现象大量TIME_WAIT状态的连接解决启用SO_REUSEADDR选项import socket sock socket.socket(socket.AF_INET, socket.SOCK_STREAM) sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)内存泄漏多进程方案中注意及时关闭Manager协程中确保正确await所有任务异常处理为每个worker设置独立异常捕获使用sentinel对象标记失败请求6. 混合模式探索对于复杂爬虫场景我们可以组合多种技术6.1 进程协程架构async def async_worker(url_queue, result_queue): async with AsyncCrawler() as crawler: while True: url await url_queue.get() if url is None: # 终止信号 break _, result await crawler.crawl(url, 1) await result_queue.put(result[0]) def process_main(urls): loop asyncio.new_event_loop() asyncio.set_event_loop(loop) url_queue asyncio.Queue() result_queue asyncio.Queue() workers [ loop.create_task(async_worker(url_queue, result_queue)) for _ in range(10) ] # 填充任务队列 for url in urls: loop.run_until_complete(url_queue.put(url)) # 等待完成 loop.run_until_complete(url_queue.join()) # 清理worker for _ in workers: loop.run_until_complete(url_queue.put(None)) loop.run_until_complete(asyncio.wait(workers)) # 收集结果 results [] while not result_queue.empty(): results.append(loop.run_until_complete(result_queue.get())) loop.close() return results if __name__ __main__: with Pool(4) as pool: chunk_size len(urls) // 4 chunks [urls[i:ichunk_size] for i in range(0, len(urls), chunk_size)] pool.map(process_main, chunks)这种架构的特点外层多进程利用多核CPU每个进程内运行事件循环驱动多个协程通过队列实现进程间任务分配6.2 性能实测对比测试10万个URL抓取任务架构模式总耗时CPU利用率内存峰值纯多进程(16进程)142s95%4.2GB纯协程(1000并发)98s65%1.1GB混合(4进程×250协程)76s88%2.3GB混合模式展现出最佳性价比在CPU利用率和内存消耗间取得了良好平衡。
返回列表