You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何改进异步爬虫架构,优化代理管理与并发任务控制?

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,避免因任务过多导致封禁。


三、递归重试的合理性

递归重试不合理,存在两个核心问题:

  1. 栈溢出风险:如果重试次数过多,递归调用会不断堆积调用栈,最终导致栈溢出。
  2. 任务泄漏:原代码中递归时未使用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

四、架构改进与优化点

  1. 代理池管理

    • 实现代理健康检测:定期向测试URL发送请求,剔除失效代理,自动补充新代理。
    • 代理使用限制:给每个代理设置请求频率上限,避免单个代理短时间内请求过多被封禁。
  2. 任务队列优化

    • 用asyncio.Queue实现生产者-消费者模型:生产者将待爬URL放入队列,消费者从队列取URL处理,结合Semaphore控制并发,支持动态添加任务。
  3. 持久化优化

    • 使用aiofiles异步文件操作库,避免同步IO阻塞事件循环。
    • 批量写入结果:积累一定数量的结果后再写入文件,减少IO操作次数。
  4. 异常与去重优化

    • 细分异常类型:针对超时、连接错误、4xx/5xx状态码分别处理,比如403直接放弃重试,5xx则更换代理重试。
    • 高效去重:用集合存储已爬取的URL/文件名,替代os.listdir的线性查询;或用数据库(如SQLite)记录爬取状态。
  5. 请求模拟优化

    • 随机切换User-Agent、Referer等请求头,模拟真实浏览器行为。
    • 添加请求间隔:即使控制并发,也给每个请求添加随机短间隔,避免过于密集。
  6. 日志系统

    • 用logging模块替代print,记录请求状态、错误信息、代理使用情况等,方便调试和问题排查。

内容的提问来源于stack exchange,提问作者sushiwithoutsushi

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 06:05:57