Python并发爬虫实战:多线程与多进程架构设计与性能优化

Python并发爬虫实战:多线程与多进程架构设计与性能优化
1. 项目概述为什么需要并发爬虫做爬虫的朋友尤其是从单线程脚本起步的几乎都经历过一个痛苦的阶段面对成百上千个页面脚本吭哧吭哧地跑一个页面卡住整个程序就跟着“罚站”效率低得让人抓狂。我刚开始写爬虫时一个简单的商品列表抓取跑完一晚上第二天一看才处理了不到十分之一那种挫败感记忆犹新。问题的核心就在于网络请求是典型的I/O密集型操作——大部分时间都在等待服务器的响应CPU却在悠闲地“喝茶”。单线程模型下这种等待是阻塞的造成了巨大的资源浪费。于是并发技术就成了爬虫效率提升的“救命稻草”。在Python的世界里我们主要有多线程threading和多进程multiprocessing两把利器。简单来说多线程就像在一个厨房里让一个厨师同时照看几个灶台线程哪个锅里的菜好了I/O完成就去处理哪个适合I/O密集的任务。而多进程则是直接开了几个独立的厨房进程每个厨房都有自己的厨师和全套工具独立内存空间适合计算密集或需要规避GIL全局解释器锁限制的任务。对于爬虫这种绝大部分时间都在进行网络I/O的场景多线程通常是首选因为它更轻量创建和切换开销小。但当任务量极大或者需要规避GIL对纯Python代码执行效率的影响时多进程方案就派上用场了。这个项目就是带你从零开始一步步把一个慢吞吞的单线程爬虫改造为能同时处理数十、上百个请求的高效并发爬虫。我们会深入探讨线程与进程的选择策略、任务队列的设计、资源竞争的处理以及如何优雅地处理异常和关闭程序。无论你是刚接触并发编程的新手还是想优化现有爬虫的老手相信都能从中获得实用的思路和可直接复用的代码。2. 核心思路与架构设计在动手写代码之前理清思路和设计好架构至关重要。一个鲁棒的并发爬虫绝不是简单地把requests.get()扔进线程池就完事了。我们需要考虑任务如何分发、数据如何收集、异常如何捕获、资源如何限制等一系列问题。2.1 并发模型选型线程 vs 进程这是第一个需要做出的关键决策。选择依据主要看任务类型和资源环境。多线程 (threading)适用场景I/O密集型任务如网络请求、文件读写。爬虫的绝大多数时间花在等待服务器返回数据上这正是多线程的用武之地。优势创建和切换开销极小共享内存使得线程间通信如下文要讲的任务队列非常方便高效。劣势受Python的GIL全局解释器锁限制。对于纯CPU计算任务多线程无法实现真正的并行因为同一时刻只有一个线程能执行Python字节码。但对于爬虫GIL的影响微乎其微因为线程在等待I/O时会释放GIL。结论对于绝大多数网页抓取、API调用类爬虫多线程是首选方案。多进程 (multiprocessing)适用场景CPU密集型任务或需要完全隔离执行环境、规避GIL的场景。例如抓取到的数据需要立即进行复杂的文本分析、图像处理。优势每个进程拥有独立的Python解释器和内存空间能实现真正的多核并行计算不受GIL束缚。劣势创建和切换开销大进程间通信IPC比线程间通信复杂且慢。结论当爬虫任务中混杂了大量本地计算或者你需要抓取的网站反爬策略极其严格每个任务需要完全独立的“指纹”如独立代理IP、独立浏览器环境时考虑使用多进程。注意对于超大规模爬虫例如需要数万个并发连接可以考虑异步IOasyncioaiohttp它在单线程内通过事件循环处理海量并发连接资源利用率更高。但学习曲线较陡且代码风格与同步编程差异较大。本项目聚焦于更通用、更容易理解和调试的线程/进程模型。2.2 生产者-消费者模式任务队列的核心无论选择线程还是进程一个高效的任务分发机制是核心。最经典的模式就是生产者-消费者模式。生产者负责生成待抓取的URL。可以是一个单独的线程/进程从文件、数据库或种子列表中读取URL并放入一个共享的任务队列。任务队列一个线程安全的容器用于存放待处理的任务URL。Python的queue.Queue用于多线程和multiprocessing.Queue用于多进程是现成的解决方案。消费者即我们的工作线程/进程。它们不断地从任务队列中获取URL执行抓取、解析、存储等操作。这种模式解耦了任务生成和任务执行使得我们可以轻松控制工作线程/进程的数量实现负载均衡。队列为空时消费者会自动等待队列满时如果设置了大小生产者会自动等待天然实现了流量控制。2.3 整体架构流程图文字描述一个典型的并发爬虫工作流如下初始化创建任务队列、结果队列可选、线程池/进程池。启动生产者将初始的种子URL列表放入任务队列。对于需要翻页或链式抓取的爬虫生产者可能也是一个常驻线程负责从已抓取的内容中提取新URL并放入队列。启动消费者创建并启动N个工作线程/进程。每个消费者循环执行从任务队列获取URL - 发送HTTP请求 - 解析响应 - 处理数据保存到文件/数据库- 可能提取新URL并放回任务队列 - 标记任务完成。协调与终止等待任务队列被完全处理空然后优雅地关闭所有工作单元。通常使用“毒丸”Poison Pill信号或设置超时等待来实现。结果汇总从结果队列中收集所有处理后的数据进行最终的处理或保存。3. 多线程爬虫实战从零构建我们以一个抓取某个公开图书网站列表页和详情页的爬虫为例演示如何用多线程实现。3.1 基础工具与环境准备首先确保你的环境已经安装了必要的库。我们主要使用requests进行HTTP请求BeautifulSoup4进行HTML解析当然还有Python标准库的threading和queue。pip install requests beautifulsoup4为了模拟真实场景我们还需要处理一些反爬策略比如设置用户代理User-Agent和请求间隔。import requests from bs4 import BeautifulSoup import time import random from threading import Thread, Lock from queue import Queue import json # 定义一个简单的请求头模拟浏览器 HEADERS { User-Agent: Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/91.0.4472.124 Safari/537.36 } # 一个简单的延时函数避免请求过快 def random_delay(min_s1, max_s3): time.sleep(random.uniform(min_s, max_s))3.2 实现线程安全的任务队列与工作者这是多线程爬虫的心脏部分。我们将创建一个Worker类每个Worker实例都是一个独立的工作线程。class CrawlerWorker(Thread): 爬虫工作线程 def __init__(self, task_queue, result_queue, error_queue, lock, worker_id): super().__init__() self.task_queue task_queue # 任务队列 (URL) self.result_queue result_queue # 结果队列 (抓取到的数据) self.error_queue error_queue # 错误队列 (失败的URL和原因) self.lock lock # 线程锁用于控制台输出等同步操作 self.worker_id worker_id self.session requests.Session() # 使用Session保持连接提升效率 self.session.headers.update(HEADERS) def run(self): 线程的主循环 while True: try: # 从任务队列获取URLblockTrue表示队列空时线程会等待 # timeout5 是安全措施防止线程永远阻塞在空队列上 url self.task_queue.get(timeout5) except Exception: # 主要是queue.Empty异常 # 如果超过5秒还没拿到任务认为所有任务已处理完线程退出 with self.lock: print(fWorker-{self.worker_id}: 任务队列长时间为空线程退出。) break # 如果收到“毒丸”信号例如None也退出 if url is None: self.task_queue.task_done() break # 执行真正的抓取任务 self.fetch_page(url) # 非常重要告诉队列这个任务已经处理完成 self.task_queue.task_done() def fetch_page(self, url): 抓取单个页面 try: with self.lock: print(fWorker-{self.worker_id} 正在抓取: {url}) response self.session.get(url, timeout10) response.raise_for_status() # 如果状态码不是200抛出HTTPError # 根据URL类型进行不同的解析处理 if /list/ in url: data self.parse_list_page(response.text, url) # 将解析出的详情页URL作为新任务放入队列 for detail_url in data.get(detail_urls, []): self.task_queue.put(detail_url) # 列表页本身的数据可能也需要保存 if data.get(page_info): self.result_queue.put((list, data[page_info])) elif /detail/ in url: data self.parse_detail_page(response.text, url) self.result_queue.put((detail, data)) else: # 其他类型的页面处理... pass # 礼貌性延迟避免对服务器造成压力 random_delay(0.5, 1.5) except requests.exceptions.RequestException as e: with self.lock: print(fWorker-{self.worker_id} 抓取失败 [{url}]: {e}) self.error_queue.put((url, str(e))) except Exception as e: with self.lock: print(fWorker-{self.worker_id} 解析异常 [{url}]: {e}) self.error_queue.put((url, fParse Error: {e})) def parse_list_page(self, html, url): 解析列表页提取详情页链接和翻页信息 soup BeautifulSoup(html, html.parser) data {detail_urls: [], page_info: {}} # 假设详情页链接在 classbook-item 的a标签里 for item in soup.select(.book-item a): href item.get(href) if href and href.startswith(/detail/): # 拼接完整URL这里需要根据目标网站结构调整 full_url https://example.com href data[detail_urls].append(full_url) # 这里可以添加解析翻页的逻辑 data[page_info][url] url data[page_info][count] len(data[detail_urls]) return data def parse_detail_page(self, html, url): 解析详情页提取图书信息 soup BeautifulSoup(html, html.parser) book_info { url: url, title: soup.select_one(h1.title).get_text(stripTrue) if soup.select_one(h1.title) else , author: soup.select_one(.author).get_text(stripTrue) if soup.select_one(.author) else , price: soup.select_one(.price).get_text(stripTrue) if soup.select_one(.price) else , # ... 其他字段 } return book_info3.3 主控程序组装与调度现在我们需要一个主程序来初始化队列、创建工作线程、投放种子任务并等待所有工作完成。def main(): # 初始化队列 task_queue Queue() # 任务队列 result_queue Queue() # 结果队列 error_queue Queue() # 错误队列 console_lock Lock() # 控制台打印锁避免输出混乱 # 种子URL seed_urls [ https://example.com/books/list/1, https://example.com/books/list/2, # ... 更多初始列表页 ] # 将种子URL放入任务队列 for url in seed_urls: task_queue.put(url) # 创建工作线程池 worker_threads [] num_workers 5 # 线程数量根据网络和机器情况调整 for i in range(num_workers): worker CrawlerWorker(task_queue, result_queue, error_queue, console_lock, i) worker.start() worker_threads.append(worker) # 等待所有任务被处理完成 (阻塞主线程) task_queue.join() # 任务完成后向每个工作线程发送“毒丸”信号通知其退出 for _ in range(num_workers): task_queue.put(None) # 等待所有工作线程结束 for worker in worker_threads: worker.join() print(\n 所有抓取任务完成 \n) # 处理结果 all_results {list: [], detail: []} while not result_queue.empty(): data_type, data result_queue.get() all_results[data_type].append(data) # 保存结果到文件 with open(book_details.json, w, encodingutf-8) as f: json.dump(all_results[detail], f, ensure_asciiFalse, indent2) print(f已保存 {len(all_results[detail])} 条图书详情到 book_details.json) # 输出错误信息 if not error_queue.empty(): print(\n 抓取错误汇总 ) errors [] while not error_queue.empty(): errors.append(error_queue.get()) with open(crawl_errors.log, w, encodingutf-8) as f: for url, reason in errors: f.write(f{url}\t{reason}\n) print(f共有 {len(errors)} 个错误已记录到 crawl_errors.log) if __name__ __main__: main()这个主程序清晰地展示了控制流准备 - 启动 - 等待 - 清理 - 汇总。task_queue.join()是关键它会阻塞主线程直到队列中所有任务都被标记为task_done()。4. 进阶多进程爬虫的实现与考量当我们的任务需要更强的隔离性或者涉及大量本地计算时就需要考虑多进程。Python的multiprocessing模块提供了与threading非常相似的接口但底层是进程。4.1 多进程与多线程的关键差异内存隔离进程间不共享内存。这意味着全局变量在一个进程中的修改不会影响到另一个进程。数据通信必须通过Queue、Pipe或Manager等IPC机制。GIL规避每个进程有独立的Python解释器因此可以充分利用多核CPU进行并行计算。启动开销进程的创建和销毁比线程慢得多资源占用也更高。代码调整multiprocessing.Queue代替queue.Queuemultiprocessing.Process代替threading.Thread并且在Windows系统上多进程代码必须放在if __name__ __main__:保护块内执行。4.2 多进程爬虫代码改造要点我们将上面的多线程爬虫改造为多进程版本主要变化如下import requests from bs4 import BeautifulSoup import time, random, json # 关键变化引入 multiprocessing from multiprocessing import Process, Queue, Lock as PLock, cpu_count class ProcessCrawlerWorker(Process): # 继承 Process def __init__(self, task_queue, result_queue, error_queue, lock, worker_id): super().__init__() self.task_queue task_queue # multiprocessing.Queue self.result_queue result_queue self.error_queue error_queue self.lock lock # multiprocessing.Lock self.worker_id worker_id # 注意Session不能作为成员变量在进程间共享需要在run方法内创建 # self.session requests.Session() # 错误 def run(self): # 每个进程内部创建自己的Session session requests.Session() session.headers.update(HEADERS) while True: try: url self.task_queue.get(timeout10) # 超时可以设长一点 except Exception: with self.lock: print(fProcessWorker-{self.worker_id}: 任务队列空退出。) break if url is None: self.task_queue.task_done() break self.fetch_page(url, session) # 传入session self.task_queue.task_done() def fetch_page(self, url, session): # 增加session参数 try: with self.lock: print(fProcessWorker-{self.worker_id} 抓取: {url}) response session.get(url, timeout15) # 使用传入的session response.raise_for_status() # ... 解析逻辑与线程版相同 ... # 假设解析出数据 parsed_data self.result_queue.put(parsed_data) random_delay(1, 2) except Exception as e: with self.lock: print(fProcessWorker-{self.worker_id} 失败 [{url}]: {e}) self.error_queue.put((url, str(e))) def main_process(): # 使用 multiprocessing 的队列和锁 task_queue Queue() result_queue Queue() error_queue Queue() console_lock PLock() seed_urls [https://example.com/books/list/1, ...] for url in seed_urls: task_queue.put(url) worker_processes [] # 进程数通常设置为CPU核心数或略多对于I/O密集型可以更多 num_workers cpu_count() * 2 for i in range(num_workers): worker ProcessCrawlerWorker(task_queue, result_queue, error_queue, console_lock, i) worker.start() worker_processes.append(worker) # 注意multiprocessing.Queue 没有 join() 方法。 # 我们需要其他方式等待任务完成比如用一个计数器。 # 这里采用一个简单但有效的方法等待任务队列被取空通过消费者发送结束信号后判断 # 更健壮的做法是使用 multiprocessing.JoinableQueue 和 task_done/join # 但为简化我们使用标志位和循环检查。 # 首先等待所有初始任务被放入队列的进程结束生产者结束 # 本例中生产者是主进程投放完种子任务就结束了。 # 然后我们需要等待任务队列变空。但队列空不代表消费者处理完了因为消费者可能刚取走任务。 # 更标准的做法是使用 JoinableQueue from multiprocessing import JoinableQueue # task_queue JoinableQueue() # 改用 JoinableQueue # 在 worker 中 self.task_queue.task_done() # 在主程序中 task_queue.join() # 这会阻塞直到所有任务被标记完成 # 由于代码结构已定我们采用一个简单的活跃进程检查循环 print(主进程等待工作进程结束...) for worker in worker_processes: worker.join() # 等待每个进程结束 print(所有工作进程已结束。) # 收集结果 all_details [] while not result_queue.empty(): all_details.append(result_queue.get()) with open(book_details_mp.json, w, encodingutf-8) as f: json.dump(all_details, f, ensure_asciiFalse, indent2) print(f多进程版保存 {len(all_details)} 条数据。) if __name__ __main__: # 多进程必须有的保护 main_process()关键提示在多进程模型中requests.Session对象不能在__init__中初始化并作为成员变量共享。因为Session对象包含套接字连接等状态无法在进程间安全地序列化和传递。必须在每个进程的run方法内部创建自己的Session实例。这是多进程编程中一个非常容易踩的坑。4.3 进程池更便捷的管理方式对于固定的、可重复的任务单元使用multiprocessing.Pool进程池是更优雅的选择。它自动管理进程的生命周期和任务分配。from multiprocessing import Pool, Manager def fetch_single_page(url): 一个独立的抓取函数将被进程池中的进程调用 try: session requests.Session() session.headers.update(HEADERS) response session.get(url, timeout10) response.raise_for_status() # 解析逻辑... soup BeautifulSoup(response.text, html.parser) title soup.title.string if soup.title else No Title time.sleep(random.uniform(0.5, 1.5)) return {url: url, title: title, status: success} except Exception as e: return {url: url, error: str(e), status: fail} def main_with_pool(): urls [https://example.com/page/1, https://example.com/page/2, ...] * 10 # 很多URL # 使用Manager来创建进程间共享的列表用于收集结果对于简单场景也可以用Pool的map返回值 with Manager() as manager: results manager.list() # 共享列表 # 创建进程池进程数通常为CPU核心数 with Pool(processes4) as pool: # 使用map_async非阻塞提交任务并指定回调函数收集结果 async_result pool.map_async(fetch_single_page, urls, callbackresults.extend) # 等待所有任务完成 async_result.wait() # 进程池with块结束后会自动关闭和join所有进程 print(f所有任务完成共处理{len(results)}条结果。) # 处理结果...进程池非常适合“任务列表明确且相互独立”的批处理场景代码简洁。但对于需要动态生成任务生产者-消费者的复杂爬虫手动管理Process和Queue的灵活性更高。5. 性能调优与高级技巧构建出能跑的并发爬虫只是第一步让它跑得又快又稳才是挑战。5.1 关键参数调优并发数线程/进程数这不是越大越好。线程数受限于网络带宽和目标服务器承受能力。一般从10-20开始测试逐步增加观察机器网络使用率和目标网站响应情况。如果出现大量超时或连接被拒说明并发过高。进程数受限于CPU核心数和内存。通常设置为CPU核心数或核心数的1-2倍。过多的进程会导致大量上下文切换开销反而降低性能。测试方法写一个简单的测试脚本用不同并发数抓取固定数量的页面记录总耗时。你会找到一个“拐点”超过这个点后增加并发数带来的收益微乎其微甚至为负。超时设置requests.get(timeout(connect_timeout, read_timeout))。连接超时建议3-5秒读取超时建议10-30秒。必须设置否则僵死线程/进程会耗尽资源。重试机制网络请求失败是常态。实现一个简单的指数退避重试逻辑能极大提升鲁棒性。def fetch_with_retry(session, url, max_retries3): for attempt in range(max_retries): try: response session.get(url, timeout10) response.raise_for_status() return response except requests.exceptions.RequestException as e: if attempt max_retries - 1: raise e # 最后一次重试失败抛出异常 wait_time 2 ** attempt random.random() # 指数退避加随机抖动 time.sleep(wait_time) return None # 理论上不会执行到这里5.2 资源管理与反反爬策略连接复用与Session务必使用requests.Session()。它会自动保持TCP连接为同一主机发送多个请求时可以避免重复的三次握手显著提升速度。控制请求频率固定延迟time.sleep()是最简单的方法但效率低。随机延迟time.sleep(random.uniform(a, b))模拟人类行为更友好。自适应速率限制监控请求成功率/失败率动态调整请求间隔。失败率升高时自动延长等待时间。代理IP池应对IP封锁的终极武器。维护一个代理IP列表在请求时随机选取。需要定期检测代理IP的有效性。class ProxyPool: def __init__(self, proxy_list): self.proxies proxy_list self.lock threading.Lock() def get_proxy(self): with self.lock: return random.choice(self.proxies) if self.proxies else None # 在worker的fetch_page中 proxy proxy_pool.get_proxy() if proxy: response session.get(url, proxies{http: proxy, https: proxy}, timeout10)请求头随机化除了User-Agent还可以随机化Accept,Accept-Language,Referer等头部使请求看起来更像来自不同的浏览器。5.3 错误处理与日志记录一个工业级爬虫必须有完善的错误处理和日志。分级日志使用logging模块区分DEBUG、INFO、WARNING、ERROR等级别。将调试信息、常规操作、警告和错误分别记录便于排查。import logging logging.basicConfig(levellogging.INFO, format%(asctime)s - %(name)s - %(levelname)s - %(message)s, handlers[logging.FileHandler(crawler.log), logging.StreamHandler()]) logger logging.getLogger(__name__) # 在代码中使用 logger.info(fWorker-{self.id} started.) logger.error(fFailed to fetch {url}: {e}, exc_infoTrue)异常细分捕获具体的异常如requests.exceptions.ConnectionError,requests.exceptions.Timeout,requests.exceptions.HTTPError并采取不同的处理策略如立即重试、更换代理、将URL放入重试队列等。状态持久化对于长时间运行的爬虫需要定期将任务队列、已抓取URL集合等状态保存到文件或数据库。这样在程序崩溃重启后可以从中断处继续避免重复抓取。6. 常见问题与实战排坑记录在实际开发中你会遇到各种各样奇怪的问题。下面是我踩过的一些坑和解决方案。6.1 线程/进程卡死或无响应症状程序运行一段时间后CPU占用率很低日志停止输出但程序不退出。可能原因与排查队列阻塞生产者生产任务的速度远低于消费者处理速度导致队列满生产者put操作阻塞。或者消费者在get时未设置超时而队列一直为空。解决为queue.put()和queue.get()设置合理的timeout参数并在超时后做相应处理如记录日志、退出线程。使用Queue(maxsize)限制队列大小。网络请求僵死requests请求没有设置超时或者超时时间过长遇到一个永不响应的服务器线程就会一直挂起。解决务必为所有网络请求设置超时。timeout(3, 30)是个不错的起点。死锁多个线程竞争多个锁时如果获取锁的顺序不一致可能导致死锁。这在爬虫中相对少见但如果你自己实现了复杂的同步逻辑需要注意。解决使用threading.Lock或multiprocessing.Lock时尽量保持锁的粒度小持有时间短并确保获取和释放锁的配对。6.2 数据错乱或丢失症状保存到文件或数据库的数据不完整或者多条数据混杂在一起。可能原因与排查文件写入竞争多个线程/进程同时向同一个文件写入导致内容相互覆盖或交错。解决不要在多线程/进程中直接写同一个文件。正确的做法是让每个工作者将数据放入一个结果队列由一个专门的写入线程/进程负责从队列中取出数据并写入文件或数据库。这样是线程/进程安全的。全局变量污染在多线程中修改共享的全局变量如列表、字典而没有加锁。解决使用线程安全的queue.Queue进行数据传递或者在对共享变量进行操作时使用threading.Lock进行保护。6.3 内存泄漏或占用过高症状程序运行时间越长占用的内存越多。可能原因与排查请求响应体未释放response.content或response.text如果被长期引用例如添加到某个全局列表会导致大量的HTML/JSON数据常驻内存。解决尽快解析响应提取出需要的有用数据如标题、价格等结构化信息然后丢弃原始的响应文本。确保对响应对象的引用及时被垃圾回收。队列堆积任务队列或结果队列中的数据没有被及时消费导致大量数据堆积在内存中。解决优化消费者处理速度或者限制队列的最大大小(maxsize)。监控队列长度如果持续增长需要报警。多进程的Manager对象multiprocessing.Manager创建的共享对象如Manager().list()通信开销较大如果频繁操作大数据内存和性能开销会剧增。解决考虑使用更高效的IPC方式如multiprocessing.Queue或者让每个进程将结果写入独立的文件最后由主进程合并。6.4 被目标网站封禁症状请求开始大量返回403、429状态码或者需要验证码。应对策略降低频率这是首要措施。大幅增加请求间隔加入随机抖动。使用代理IP搭建或购买可靠的代理IP池并实现自动切换。完善请求头模拟真实浏览器的完整请求头包括Accept,Accept-Encoding,Accept-Language,Referer等。使用Cookies对于一些需要登录或会话的网站使用Session对象自动管理Cookies。设置请求间隔更高级的做法是为每个目标域名设置独立的请求间隔计时器确保遵守网站的robots.txt规则如果有的话。识别验证码如果遇到验证码需要考虑接入打码平台或者使用机器学习库进行简单识别难度较高。最后也是最重要的心得并发爬虫的调试比单线程复杂得多。出现问题时先尝试将并发数设为1回归单线程模式看问题是否复现。如果单线程正常那么多半是并发控制逻辑锁、队列出了问题。善用日志在每个关键步骤如获取任务、开始请求、收到响应、保存数据都打上带线程/进程ID的日志这样在排查问题时能清晰地看到每个工作单元的执行轨迹。