Python数据库驱动爬虫如何通过multiprocess/multithreading提速
并发爬取改造方案
网页爬取属于典型IO密集型场景,绝大多数耗时都在等待网络响应,CPU占用极低,完全没必要上多进程——多进程的启动、上下文切换开销远高于收益,还会额外带来数据库连接复用的麻烦,优先选多线程或者异步协程就够,改造成本低,性能提升明显。
方案1:线程池实现(改动最小,兼容原有逻辑)
这个方案几乎不用改你原来写的scrape_short_webinfo逻辑,用标准库自带的concurrent.futures.ThreadPoolExecutor就行,上手最快。
注意几个必做的优化点,不然跑起来容易出问题:
- 不要一次性把70万条URL全加载到内存,用offset+limit分批从数据库读,每批处理1000-5000条即可
- 线程数不要盲目开高,根据带宽和目标站反爬强度设20-50个足够,开太高反而容易占满带宽、触发反爬
- 爬取逻辑必须加超时、异常捕获,避免单个慢请求/异常直接卡死整个工作线程
参考实现代码:
import requests from concurrent.futures import ThreadPoolExecutor, as_completed # 替换成你自己项目的model导入路径 from myapp.models import Link def scrape_short_webinfo(url): # 保留你原来的爬取逻辑,只需要补全异常捕获和超时设置 try: # 示例请求,timeout必须加,单位秒 resp = requests.get(url, timeout=10, headers={"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36"}) resp.raise_for_status() # 这里写你提取短信息的逻辑 page_data = resp.text[:200] # 可以直接在这里更新数据库,或者把结果返回后批量更新 return {"url": url, "data": page_data, "code": 0} except Exception as e: return {"url": url, "data": None, "code": -1, "msg": str(e)} def batch_get_urls(batch_size=2000): # 分批读URL,避免一次性加载全量数据占满内存 offset = 0 while True: url_batch = list( Link.objects.all() .order_by("id")[offset:offset+batch_size] .values_list("url", flat=True) ) if not url_batch: break yield url_batch offset += batch_size if __name__ == "__main__": # 线程数按需调整 WORKER_NUM = 30 for urls in batch_get_urls(): with ThreadPoolExecutor(max_workers=WORKER_NUM) as executor: task_map = {executor.submit(scrape_short_webinfo, url): url for url in urls} for task in as_completed(task_map): result = task.result() # 处理结果:打日志、更新数据库都可以,建议攒一批批量更新库,减少IO print(f"处理完成: {result['url']}, 状态: {result['code']}")
方案2:异步协程实现(性能更高)
如果线程池的速度还是达不到要求,可以换asyncio+aiohttp的异步协程方案,单线程就能支撑上百并发,资源开销比线程池更低,速度更快。缺点是原来的同步爬取逻辑需要改成异步写法,改造成本稍高。
参考实现:
import asyncio import aiohttp from myapp.models import Link async def scrape_task(session, url, sem): # sem用来控制全局并发数,避免请求发太猛 async with sem: try: async with session.get( url, timeout=aiohttp.ClientTimeout(total=10), headers={"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36"} ) as resp: text = await resp.text(errors="ignore") return {"url": url, "data": text[:200], "code": 0} except Exception as e: return {"url": url, "data": None, "code": -1, "msg": str(e)} async def main(): # 最大并发数,按需调整 MAX_CONCURRENT = 50 sem = asyncio.Semaphore(MAX_CONCURRENT) batch_size = 2000 offset = 0 # 连接器关闭默认并发限制,交给sem统一控制,开DNS缓存减少重复解析开销 connector = aiohttp.TCPConnector(limit=0, ttl_dns_cache=300) async with aiohttp.ClientSession(connector=connector) as session: while True: url_batch = list( Link.objects.all() .order_by("id")[offset:offset+batch_size] .values_list("url", flat=True) ) if not url_batch: break tasks = [scrape_task(session, url, sem) for url in url_batch] for res in asyncio.as_completed(tasks): result = await res print(f"处理完成: {result['url']}, 状态: {result['code']}") offset += batch_size if __name__ == "__main__": asyncio.run(main())
跑之前必看的避坑点
- 不要用多进程:爬取是IO密集型场景,多进程的额外开销会抵消并发收益,还需要处理跨进程的数据库连接、结果传递问题,性价比极低
- 数据库写入尽量批量做:不要每爬完一条就单独update一次数据库,攒100-200条结果用
bulk_update一次性写入,能大幅降低数据库压力 - 并发数不是越大越好:普通家用带宽开20-30并发就够,服务器带宽充足的话最多开到100-200,太高会导致所有请求因为带宽不足超时,还容易触发目标站反爬封IP
- 记得加随机UA、必要时上代理池:如果爬的站点反爬严格,固定UA爬不了多少就会被拦截,需要定期切换请求身份
- 一定要给请求设超时:不管用哪种方案,超时参数必须加,不然遇到死链、慢响应的站点,工作线程/协程会被长期占用,最后整个程序卡死
内容的提问来源于stack exchange,提问作者user14880125
相关产品推荐
相关产品推荐

