Python多进程爬虫/数据处理实战:用apply_async和回调函数优雅处理异常与结果
·
Python多进程爬虫与数据处理实战:异常处理与结果回调的工程化实践
当面对海量数据抓取或复杂计算任务时,单进程处理往往成为性能瓶颈。我曾接手过一个电商价格监控项目,最初使用单线程爬虫每小时仅能处理200个商品页面,直到引入multiprocessing.Pool的apply_async方法配合回调机制,才将效率提升到每小时8000+页面,同时保证了系统的稳定性。本文将分享如何构建生产级的多进程处理框架,特别是异常处理和结果收集的工程实践。
1. 多进程任务的核心架构设计
在构建健壮的多进程系统前,需要理解几个关键设计原则。进程池的大小通常设置为CPU核心数的1-5倍,I/O密集型任务可以适当增加,而计算密集型任务则建议保持较小规模。以下是不同场景下的进程池配置建议:
| 任务类型 | CPU核心数 | 推荐进程数 | 内存消耗 |
|---|---|---|---|
| 纯计算任务 | 4 | 4 | 低 |
| 网络请求密集型 | 4 | 8-12 | 中 |
| 混合型任务 | 4 | 6-8 | 中高 |
典型的任务处理流程应该包含以下组件:
from multiprocessing import Pool
import traceback
def worker(task_data):
try:
# 核心业务逻辑
result = process_data(task_data)
return {'status': 'success', 'data': result}
except Exception as e:
return {'status': 'error', 'exception': str(e), 'traceback': traceback.format_exc()}
def success_callback(result):
if result['status'] == 'success':
store_result(result['data'])
else:
handle_failure(result)
def error_callback(error):
log_error(f"进程崩溃: {str(error)}")
pool = Pool(processes=8)
for task in task_generator():
pool.apply_async(
worker,
args=(task,),
callback=success_callback,
error_callback=error_callback
)
2. 异常处理的工程级解决方案
在实际项目中,简单的try-catch往往不足以应对复杂的异常场景。我们需要建立分层的异常处理机制:
- 基础异常捕获:处理已知的常见错误(网络超时、数据格式异常等)
- 进程级防护:防止单个进程崩溃影响整个任务队列
- 资源回收:确保异常发生时正确释放资源
一个电商爬虫的异常处理示例:
def fetch_product_page(url):
try:
response = requests.get(url, timeout=10)
response.raise_for_status()
return parse_html(response.text)
except requests.exceptions.RequestException as e:
raise PageFetchError(f"页面获取失败: {url}") from e
except HTMLParseError as e:
raise DataParseError(f"HTML解析错误: {url}") from e
def error_handler(error):
error_type = type(error).__name__
if error_type in ('PageFetchError', 'DataParseError'):
logger.warning(f"可恢复错误: {str(error)}")
retry_queue.put(error.args[0]) # 加入重试队列
else:
logger.error(f"严重错误: {traceback.format_exc()}")
emergency_stop() # 触发紧急停止
常见异常处理策略对比:
| 策略类型 | 适用场景 | 实现复杂度 | 资源消耗 |
|---|---|---|---|
| 直接丢弃 | 非关键任务 | 低 | 低 |
| 立即重试 | 临时性错误 | 中 | 中 |
| 加入重试队列 | 可延迟任务 | 高 | 中 |
| 全局终止 | 致命错误 | 中 | 高 |
3. 高级回调模式与结果处理
当处理数百万级任务时,简单的结果收集方式会导致内存爆炸。以下是几种经过实战检验的模式:
流式处理回调:
class ResultCollector:
def __init__(self):
self.counter = 0
self.buffer = []
def add_result(self, result):
self.buffer.append(result)
if len(self.buffer) >= 1000:
self.flush()
def flush(self):
save_to_database(self.buffer)
self.counter += len(self.buffer)
self.buffer = []
collector = ResultCollector()
def callback(result):
collector.add_result(result)
if collector.counter % 10000 == 0:
logger.info(f"已处理 {collector.counter} 条记录")
分布式结果队列(适用于跨机器场景):
import redis
r = redis.Redis(host='redis-host')
def distributed_callback(result):
try:
r.rpush('result_queue', json.dumps(result))
except Exception as e:
logger.error(f"Redis写入失败: {e}")
local_fallback_storage.save(result) # 降级方案
4. 性能优化与实战技巧
经过多个项目的性能分析,我们发现多进程系统的瓶颈往往出现在以下几个方面:
- 进程间通信开销:特别是当传递大数据量时
- 全局锁竞争:如同时写入同一个文件
- 资源争用:数据库连接、网络带宽等
优化后的文件处理模板:
def process_file_chunk(args):
file_path, offset, size = args
results = []
with open(file_path, 'rb') as f:
f.seek(offset)
chunk = f.read(size)
for line in chunk.splitlines():
try:
results.append(parse_line(line))
except Exception as e:
results.append({'error': str(e), 'line': line})
return results
def chunked_file_processor(file_path, chunk_size=1024*1024):
file_size = os.path.getsize(file_path)
chunks = [(file_path, i, min(chunk_size, file_size-i))
for i in range(0, file_size, chunk_size)]
with Pool(processes=4) as pool:
for result in pool.imap_unordered(process_file_chunk, chunks):
yield from result
关键性能指标监控(使用psutil实现):
import psutil
def monitor_resources():
while True:
cpu_percent = psutil.cpu_percent(interval=1)
mem_info = psutil.virtual_memory()
if cpu_percent > 90:
throttle_processing() # 降低处理速度
if mem_info.percent > 80:
trigger_gc() # 主动触发垃圾回收
5. 复杂场景下的稳定性保障
在长时间运行的生产环境中,我们需要考虑更多边缘情况:
- 僵尸进程防护:添加进程心跳检测
- 内存泄漏预防:定期重启工作进程
- 优雅退出:处理中断信号
进程管理增强版:
from multiprocessing import Pool, Manager
import signal
class ManagedPool:
def __init__(self, size):
self.pool = Pool(size)
self.manager = Manager()
self.heartbeats = self.manager.dict()
def worker_init(self):
signal.signal(signal.SIGINT, signal.SIG_IGN)
self.heartbeats[os.getpid()] = time.time()
def check_health(self):
now = time.time()
dead = [pid for pid, ts in self.heartbeats.items()
if now - ts > 300] # 5分钟无心跳
for pid in dead:
os.kill(pid, signal.SIGTERM)
def apply_async(self, func, args=(), kwargs={}, callback=None):
def wrapped(*args, **kwargs):
self.worker_init()
return func(*args, **kwargs)
return self.pool.apply_async(wrapped, args, kwargs, callback)
在实际项目中,我发现最有效的稳定性策略是"快速失败+自动恢复"机制。当某个工作进程连续失败超过阈值时,自动将其隔离并启动新的替代进程,同时将失败任务转移到健康节点。
更多推荐


所有评论(0)