Queue的多线程爬虫和multiprocessing多进程
Queue的多线程爬虫和multiprocessing多进程
一、背景与问题
在分布式系统和高性能计算场景中,多线程与多进程是两种核心的并发模型。对于爬虫系统而言,如何高效处理海量URL的采集任务,是决定系统性能的关键。
传统单线程爬虫在处理大量请求时会遇到明显的性能瓶颈,而多线程和多进程提供了两种不同的解决方案。但两者在适用场景、资源消耗、实现复杂度等方面存在本质差异。
以一个典型的爬虫场景为例:需要处理10万条URL,每个请求平均耗时100ms。单线程处理需要10万秒(约27小时),而使用多线程或进程池可以将时间缩短至数分钟。但具体选择哪种方案,需要深入理解其技术原理。
二、基本原理
1. 线程与进程的本质差异
- 线程:共享同一进程的内存空间,通过协程切换实现并发。受GIL(全局解释器锁)限制,CPython中线程无法实现真正的并行计算
- 进程:独立的内存空间,通过进程间通信(IPC)实现协作。每个进程有独立的Python解释器实例
2. 队列的核心作用
队列(Queue)在并发系统中扮演着任务调度中枢的角色。其核心特性包括:
- 线程安全:支持多线程安全的put()和get()操作
- 阻塞机制:在队列为空时get()会阻塞,避免忙等待
- 任务分发:将任务均匀分配给工作线程/进程
3. 线程池与进程池的对比
| 特性 | 线程池(ThreadPoolExecutor) | 进程池(ProcessPoolExecutor) |
|---|---|---|
| 资源消耗 | 低(共享内存) | 高(独立内存空间) |
| 调度粒度 | 线程级(轻量) | 进程级(重量) |
| GIL影响 | 受限于GIL | 无GIL限制 |
| 内存共享 | 可共享数据结构 | 需通过IPC传递数据 |
| 适用场景 | I/O密集型任务 | CPU密集型任务 |
三、环境准备
# 安装必要库
pip install concurrent.futures requests# 导入核心模块
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
import requests
from queue import Queue
import threading
import time四、核心实现
1. 多线程爬虫实现(线程池)
def fetch_page(url, result_queue):
try:
response = requests.get(url, timeout=10)
result_queue.put((url, len(response.text)))
except Exception as e:
result_queue.put((url, str(e)))
def thread_crawler(urls):
queue = Queue()
for url in urls:
queue.put(url)
with ThreadPoolExecutor(max_workers=10) as executor:
futures = []
for _ in range(10): # 10个线程
future = executor.submit(fetch_page, queue.get(), queue)
futures.append(future)
# 等待所有任务完成
for future in futures:
future.result()关键代码解释:
- 使用
ThreadPoolExecutor创建线程池 Queue用于任务分发和结果收集- 通过
submit()提交任务,自动管理线程生命周期 result()方法获取任务结果
2. 多进程爬虫实现(进程池)
def fetch_page_process(url, result_queue):
try:
response = requests.get(url, timeout=10)
result_queue.put((url, len(response.text)))
except Exception as e:
result_queue.put((url, str(e)))
def process_crawler(urls):
queue = Queue()
for url in urls:
queue.put(url)
with ProcessPoolExecutor(max_workers=4) as executor:
futures = []
for _ in range(4): # 4个进程
future = executor.submit(fetch_page_process, queue.get(), queue)
futures.append(future)
for future in futures:
future.result()关键代码解释:
- 使用
ProcessPoolExecutor创建进程池 - 进程间通过
multiprocessing.Queue通信 - 每个进程独立运行Python解释器
- 需要特别注意进程间的数据同步
3. 混合使用线程和进程的案例
def fetch_page(url, result_queue):
try:
response = requests.get(url, timeout=10)
result_queue.put((url, len(response.text)))
except Exception as e:
result_queue.put((url, str(e)))
def hybrid_crawler(urls):
queue = Queue()
for url in urls:
queue.put(url)
with ThreadPoolExecutor(max_workers=10) as thread_pool:
# 线程池处理I/O密集型任务
thread_futures = []
for _ in range(10):
future = thread_pool.submit(fetch_page, queue.get(), queue)
thread_futures.append(future)
# 进程池处理CPU密集型任务(假设此处需要计算)
with ProcessPoolExecutor(max_workers=4) as process_pool:
process_futures = []
for _ in range(4):
# 假设此处需要计算
future = process_pool.submit(lambda: (None, None))
process_futures.append(future)
# 等待所有任务完成
for future in thread_futures + process_futures:
future.result()关键代码解释:
- 线程池处理网络请求等I/O操作
- 进程池处理需要大量计算的任务
- 通过队列进行任务分发和结果收集
- 需要特别注意线程和进程的资源分配
五、完整案例
爬虫系统完整实现
import requests
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
from queue import Queue
import threading
import time
import random
# 模拟URL列表
urls = [
"https://example.com",
"https://example.org",
"https://example.net",
"https://example.edu",
"https://example.gov"
] * 10000 # 10000个URL
# 线程安全的计数器
class SafeCounter:
def __init__(self):
self.lock = threading.Lock()
self.count = 0
def increment(self):
with self.lock:
self.count += 1
# 线程爬虫函数
def fetch_page(url, result_queue, counter):
try:
response = requests.get(url, timeout=10)
result_queue.put((url, len(response.text)))
counter.increment()
except Exception as e:
result_queue.put((url, str(e)))
# 多线程爬虫
def thread_crawler():
queue = Queue()
counter = SafeCounter()
for url in urls:
queue.put(url)
with ThreadPoolExecutor(max_workers=10) as executor:
futures = []
for _ in range(10):
future = executor.submit(fetch_page, queue.get(), queue, counter)
futures.append(future)
# 等待所有任务完成
for future in futures:
future.result()
print(f"Total fetched: {counter.count}")
# 多进程爬虫
def process_crawler():
queue = Queue()
counter = SafeCounter()
for url in urls:
queue.put(url)
with ProcessPoolExecutor(max_workers=4) as executor:
futures = []
for _ in range(4):
future = executor.submit(fetch_page, queue.get(), queue, counter)
futures.append(future)
for future in futures:
future.result()
print(f"Total fetched: {counter.count}")
# 主程序
if __name__ == "__main__":
start_time = time.time()
thread_crawler()
print(f"Thread crawler took: {time.time() - start_time:.2f} seconds")
start_time = time.time()
process_crawler()
print(f"Process crawler took: {time.time() - start_time:.2f} seconds")关键实现细节:
- 使用线程安全计数器统计成功请求
- 队列用于任务分发和结果收集
- 线程池和进程池分别处理不同类型的任务
- 添加了性能测试指标
六、源码解析
1. 线程池源码分析
ThreadPoolExecutor的submit()方法内部:
- 创建一个
Future对象 - 将任务提交到线程池的队列中
- 选择一个空闲线程执行任务
- 通过
Future对象获取结果
关键代码:
def submit(self, fn, *args, **kwargs):
if self._shutdown:
raise RuntimeError("Cannot submit new tasks to a shut down executor")
if self._max_workers == 0:
raise ValueError("Cannot submit new tasks to a executor with zero max_workers")
future = Future()
self._work_queue.put((fn, args, kwargs, future))
self._adjust_thread_count()
return future2. 进程池源码分析
ProcessPoolExecutor的submit()方法:
- 创建一个新的进程
- 通过IPC传递任务参数
- 在新进程中执行任务
- 通过共享内存返回结果
关键代码:
def submit(self, fn, *args, **kwargs):
if self._shutdown:
raise RuntimeError("Cannot submit new tasks to a shut down executor")
if self._max_workers == 0:
raise ValueError("Cannot submit new tasks to a executor with zero max_workers")
future = Future()
self._work_queue.put((fn, args, kwargs, future))
self._adjust_process_count()
return future七、进阶使用
1. 动态调整线程/进程数量
def dynamic_crawler(urls):
queue = Queue()
for url in urls:
queue.put(url)
with ThreadPoolExecutor(max_workers=10) as thread_pool:
# 动态调整线程数
thread_pool.submit(fetch_page, queue.get(), queue)
# 根据负载动态调整
while not queue.empty():
if len(thread_pool._threads) < 10:
thread_pool.submit(fetch_page, queue.get(), queue)2. 高级队列管理
class BoundedQueue:
def __init__(self, maxsize=0):
self.queue = Queue(maxsize)
def put(self, item):
if self.queue.full():
raise QueueFullError("Queue is full")
self.queue.put(item)
def get(self):
return self.queue.get()3. 异步IO混合使用
import asyncio
from concurrent.futures import ThreadPoolExecutor
async def async_crawler(urls):
loop = asyncio.get_event_loop()
with ThreadPoolExecutor() as pool:
tasks = [loop.create_task(fetch_page_async(url, pool)) for url in urls]
await asyncio.gather(*tasks)八、性能与工程实践
1. 性能调优建议
| 优化方向 | 建议措施 | 效果说明 |
|---|---|---|
| 线程数设置 | 根据CPU核心数调整(通常为CPU*2) | 提高并发处理能力 |
| 队列大小 | 设置合理上限(如1000) | 避免内存溢出 |
| 网络超时 | 设置合理超时时间(如5秒) | 避免阻塞长时间等待 |
| 任务分片 | 将大任务拆分为小任务 | 提高资源利用率 |
| 资源回收 | 及时关闭空闲线程/进程 | 释放系统资源 |
2. 异常处理策略
- 网络异常:添加重试机制和重试策略
- 任务异常:捕获异常并记录日志
- 资源异常:设置资源限制和告警机制
3. 安全考虑
- 请求限制:设置请求频率限制,避免被封IP
- 数据验证:对返回数据进行合法性校验
- 身份认证:使用API密钥或OAuth进行身份验证
- 缓存策略:合理使用缓存减少服务器压力
九、常见问题与踩坑
1. 线程池死锁问题
错误示例:
def worker():
with lock:
do_something()问题:多个线程同时持有锁会导致死锁
解决办法:使用contextlib.contextmanager管理锁,确保锁的释放
2. 进程间通信问题
错误示例:
def worker():
print("Process started")问题:进程启动后无法立即返回结果
解决办法:使用multiprocessing.Pipe或multiprocessing.Queue进行通信
3. 资源竞争问题
错误示例:
shared_counter = 0
def increment():
shared_counter += 1问题:多线程/进程竞争导致计数错误
解决办法:使用线程锁或进程锁保护共享资源
4. 性能瓶颈问题
错误示例:未限制线程/进程数量导致资源耗尽
解决办法:根据系统资源动态调整并发数量
十、最佳实践
1. 选择原则
| 场景类型 | 推荐方案 | 理由 |
|---|---|---|
| I/O密集型 | 线程池 | 轻量级,适合网络请求 |
| CPU密集型 | 进程池 | 避免GIL限制,适合计算任务 |
| 混合场景 | 线程+进程混合 | 分工协作,发挥各自优势 |
| 高并发场景 | 异步IO+线程池 | 提高吞吐量,降低延迟 |
2. 代码规范建议
- 使用
concurrent.futures模块而非原始thread/process - 使用
queue.Queue管理任务分发 - 添加异常处理和日志记录
- 使用
with语句管理资源 - 避免直接操作线程/进程对象
3. 性能监控建议
- 使用
time模块记录任务耗时 - 使用
logging模块记录关键指标 - 使用
psutil监控系统资源 - 使用
cProfile进行性能分析
十一、总结
多线程和多进程是构建高性能爬虫系统的两大核心支柱。在实际开发中,需要根据任务类型选择合适的并发模型:I/O密集型任务适合线程池,CPU密集型任务适合进程池。通过合理使用队列管理任务分发,结合线程锁/进程锁保护共享资源,可以构建稳定高效的爬虫系统。
在实际项目中,建议遵循以下原则:
- 使用
concurrent.futures模块进行并发控制 - 通过队列管理任务分发和结果收集
- 根据系统资源动态调整并发数量
- 增加异常处理和日志记录机制
- 对关键代码进行性能测试和优化
对于大规模爬虫系统,可以结合异步IO、缓存策略、分布式架构等技术,构建更复杂的并发处理体系。同时,需要特别注意网络请求的合法性和资源限制,避免对目标服务器造成过大压力。
评论已关闭