如何改进异步爬虫架构,优化代理管理与并发任务控制?
Web爬虫asyncio与代理技术优化问题
我是一名Web爬虫爱好者,刚学习了Web爬虫中的asyncio与代理技术以提升性能。编写的程序能快速下载大量站点,但复盘后发现两大核心问题:一是运行中的任务在执行中途收到错误响应时无法更换代理;二是对并发任务数m的控制效果很差。另外还有重试方案的合理性疑问,以及希望得到架构改进和优化建议。
原代码
import asyncio import httpx import os import time async def fetch_response(client, url, t): headers = {...} resp = await client.get(url, headers=headers, timeout=30) if resp.status_code == httpx.codes.OK: return [resp.text, str(resp.url)[108:113]] else: # 响应非200时递归重试 print(f'Error: {resp.status_code}') time.sleep(n) fetch_response(client, url, t+5) async def main(): proxies = getProxyList() # 获取代理列表,当前共10个 files = os.listdir('results/') # 获取已下载的文件 # 读取URL列表 with open('urls.txt','r') as f: urls = [line.strip() for line in f.readlines()] n = 1 # 代理轮换索引 m = 10 # 并发任务数 t = 5 # 重试等待时间 results = [] # 存储任务结果 # 按批次处理URL for i in range(0,len(urls),m): async with httpx.AsyncClient(proxies={'http://':proxies[n],'https://':proxies[n]}) as client: tasks = [] for j in range(0,m): # 跳过已下载的文件对应的URL if f"result{i:0>5}.json" not in files: tasks.append(asyncio.create_task(fetch_response(client, urls[i+j], t))) results = await asyncio.gather(*tasks) # 写入结果文件 if len(results)>0: for k in results: with open(f'results/result{k[1]}.json','w+') as f: f.write(k[0]) print(i) n += 1 # 代理轮换到末尾后重置索引 if n == 10: n=0
一、代理动态切换问题
当前代码用一个代理跑完整批次m个任务后再切换,灵活性差,且无法在任务执行中途更换代理。可以实现中途切换代理,有两种可行方案:
方案1:请求时动态指定代理
无需固定AsyncClient的代理,每次请求时通过proxies参数指定代理,出错时直接换代理重新请求:
async def fetch_response(proxy_pool, url, t): headers = {...} max_retries = 3 for _ in range(max_retries): # 从代理池随机/轮换取一个代理 proxy = proxy_pool.get() try: async with httpx.AsyncClient(proxies={'http://':proxy,'https://':proxy}) as client: resp = await client.get(url, headers=headers, timeout=30) if resp.status_code == httpx.codes.OK: return [resp.text, str(resp.url)[108:113]] print(f'状态码错误: {resp.status_code}, 更换代理重试') except httpx.ProxyError: print(f'代理失效: {proxy}, 移除并更换') proxy_pool.remove(proxy) await asyncio.sleep(t) return None
方案2:复用AsyncClient但动态修改代理
httpx.AsyncClient支持修改默认代理,可通过client.proxies属性动态更新(asyncio单线程下无需担心线程安全):
async def fetch_response(client, proxy_pool, url, t): headers = {...} max_retries = 3 for _ in range(max_retries): try: resp = await client.get(url, headers=headers, timeout=30) if resp.status_code == httpx.codes.OK: return [resp.text, str(resp.url)[108:113]] print(f'状态码错误: {resp.status_code}, 更换代理') except httpx.ProxyError: print(f'代理失效, 更换') # 更换代理 client.proxies = {'http://':proxy_pool.get(),'https://':proxy_pool.get()} await asyncio.sleep(t) return None
二、并发任务数控制问题
当前按批次切片的方式无法精准控制并发数(比如批次内某任务耗时极长,会占用资源却无法启动新任务),正确做法是用asyncio.Semaphore实现精准并发控制:
async def fetch_response(sem, proxy_pool, url, t): # 用信号量限制并发数 async with sem: headers = {...} max_retries = 3 for _ in range(max_retries): proxy = proxy_pool.get() try: async with httpx.AsyncClient(proxies={'http://':proxy,'https://':proxy}) as client: resp = await client.get(url, headers=headers, timeout=30) if resp.status_code == httpx.codes.OK: return [resp.text, str(resp.url)[108:113]] except Exception as e: print(f'请求失败: {e}') await asyncio.sleep(t) return None async def main(): proxies = getProxyList() files = os.listdir('results/') with open('urls.txt','r') as f: urls = [line.strip() for line in f.readlines()] m = 10 # 最大并发数 sem = asyncio.Semaphore(m) # 过滤已爬取的URL pending_urls = [url for idx, url in enumerate(urls) if f"result{idx:0>5}.json" not in files] # 创建所有任务 tasks = [asyncio.create_task(fetch_response(sem, proxies, url, 5)) for url in pending_urls] results = await asyncio.gather(*tasks) # 批量写入结果 for idx, result in enumerate(results): if result: with open(f'results/result{idx:0>5}.json','w+') as f: f.write(result[0])
这种方式不管总任务量多大,同时运行的任务数严格限制为m,避免因任务过多导致封禁。
三、递归重试的合理性
递归重试不合理,存在两个核心问题:
- 栈溢出风险:如果重试次数过多,递归调用会不断堆积调用栈,最终导致栈溢出。
- 任务泄漏:原代码中递归时未使用
await,也未返回递归结果,会导致异步任务无法正确完成,结果丢失。
正确方案:用循环实现重试,同时替换time.sleep为asyncio.sleep(避免阻塞事件循环):
async def fetch_response(sem, proxy_pool, url, t): async with sem: headers = {...} max_retries = 3 retry_count = 0 while retry_count < max_retries: proxy = proxy_pool.get() try: async with httpx.AsyncClient(proxies={'http://':proxy,'https://':proxy}) as client: resp = await client.get(url, headers=headers, timeout=30) if resp.status_code == httpx.codes.OK: return [resp.text, str(resp.url)[108:113]] print(f'错误状态码: {resp.status_code}, 重试次数: {retry_count+1}') except httpx.TimeoutException: print(f'请求超时, 重试次数: {retry_count+1}') except httpx.ProxyError: print(f'代理失效: {proxy}, 重试次数: {retry_count+1}') retry_count += 1 await asyncio.sleep(t * retry_count) # 指数退避等待 print(f'URL {url} 重试{max_retries}次失败') return None
四、架构改进与优化点
代理池管理
- 实现代理健康检测:定期向测试URL发送请求,剔除失效代理,自动补充新代理。
- 代理使用限制:给每个代理设置请求频率上限,避免单个代理短时间内请求过多被封禁。
任务队列优化
- 用
asyncio.Queue实现生产者-消费者模型:生产者将待爬URL放入队列,消费者从队列取URL处理,结合Semaphore控制并发,支持动态添加任务。
- 用
持久化优化
- 使用
aiofiles异步文件操作库,避免同步IO阻塞事件循环。 - 批量写入结果:积累一定数量的结果后再写入文件,减少IO操作次数。
- 使用
异常与去重优化
- 细分异常类型:针对超时、连接错误、4xx/5xx状态码分别处理,比如403直接放弃重试,5xx则更换代理重试。
- 高效去重:用集合存储已爬取的URL/文件名,替代
os.listdir的线性查询;或用数据库(如SQLite)记录爬取状态。
请求模拟优化
- 随机切换User-Agent、Referer等请求头,模拟真实浏览器行为。
- 添加请求间隔:即使控制并发,也给每个请求添加随机短间隔,避免过于密集。
日志系统
- 用
logging模块替代print,记录请求状态、错误信息、代理使用情况等,方便调试和问题排查。
- 用
内容的提问来源于stack exchange,提问作者sushiwithoutsushi
相关产品推荐
相关产品推荐

