本文目录导读:

我来介绍几种Python进程池批量管理的实战案例。
基础进程池使用
简单任务分发
import multiprocessing
import time
import random
def worker_task(num):
"""模拟耗时任务"""
sleep_time = random.uniform(0.5, 2.0)
time.sleep(sleep_time)
print(f"任务 {num} 完成,耗时: {sleep_time:.2f}秒")
return num * 2
# 创建进程池
with multiprocessing.Pool(processes=4) as pool:
# 批量提交任务
results = pool.map(worker_task, range(10))
print(f"所有结果: {results}")
使用回调函数管理进度
import multiprocessing
from functools import partial
def process_file(filename, output_dir):
"""模拟文件处理"""
time.sleep(1)
result = f"{output_dir}/processed_{filename}"
return result
def progress_callback(result):
"""进度回调"""
print(f"处理完成: {result}")
# 批量管理进程
def batch_process_files(files, output_dir, num_workers=4):
with multiprocessing.Pool(num_workers) as pool:
# 使用partial固定参数
process_func = partial(process_file, output_dir=output_dir)
# 异步提交任务
async_results = []
for file in files:
result = pool.apply_async(
process_func,
args=(file,),
callback=progress_callback
)
async_results.append(result)
# 等待所有任务完成
pool.close()
pool.join()
# 收集结果
results = [r.get() for r in async_results]
return results
# 使用示例
files = [f"file_{i}.txt" for i in range(20)]
output = batch_process_files(files, "/output")
动态任务队列管理
from multiprocessing import Pool, Manager, Queue
import queue
import threading
class DynamicTaskManager:
def __init__(self, num_workers=4):
self.num_workers = num_workers
self.task_queue = Queue()
self.result_queue = Queue()
self.is_running = False
def add_task(self, task_func, *args, **kwargs):
"""添加任务到队列"""
self.task_queue.put((task_func, args, kwargs))
def worker_process(self, task_queue, result_queue, worker_id):
"""工作进程"""
while True:
try:
# 非阻塞获取任务
task_func, args, kwargs = task_queue.get(timeout=1)
result = task_func(*args, **kwargs)
result_queue.put({
'worker_id': worker_id,
'result': result,
'success': True
})
except queue.Empty:
# 队列为空,退出
break
except Exception as e:
result_queue.put({
'worker_id': worker_id,
'error': str(e),
'success': False
})
def run(self):
"""运行任务管理器"""
with Manager() as manager:
task_queue = manager.Queue()
result_queue = manager.Queue()
# 将任务放入共享队列
while not self.task_queue.empty():
task_queue.put(self.task_queue.get())
# 创建进程池
with Pool(self.num_workers) as pool:
# 启动工作进程
workers = []
for i in range(self.num_workers):
worker = pool.apply_async(
self.worker_process,
args=(task_queue, result_queue, i)
)
workers.append(worker)
# 收集结果
results = []
completed = 0
total_tasks = task_queue.qsize()
while completed < total_tasks:
try:
result = result_queue.get(timeout=2)
results.append(result)
completed += 1
print(f"进度: {completed}/{total_tasks}")
except queue.Empty:
continue
return results
# 使用示例
def custom_task(x, y):
time.sleep(0.5)
return x * y
manager = DynamicTaskManager(num_workers=3)
for i in range(10):
manager.add_task(custom_task, i, i+1)
results = manager.run()
print(f"完成 {len(results)} 个任务")
批量数据处理框架
import pandas as pd
from multiprocessing import Pool
import numpy as np
class BatchDataProcessor:
def __init__(self, chunk_size=1000, num_workers=4):
self.chunk_size = chunk_size
self.num_workers = num_workers
def process_chunk(self, chunk_data):
"""处理数据块"""
# 模拟数据处理
time.sleep(0.1)
# 示例:数据转换
if isinstance(chunk_data, pd.DataFrame):
result = chunk_data.apply(lambda x: x * 2)
else:
result = [item * 2 for item in chunk_data]
return result
def batch_process(self, data):
"""批量处理数据"""
total_items = len(data)
# 将数据分块
chunks = []
for i in range(0, total_items, self.chunk_size):
chunk = data[i:i+self.chunk_size]
chunks.append(chunk)
print(f"数据分块: {len(chunks)} 块")
# 使用进程池并行处理
with Pool(self.num_workers) as pool:
results = pool.map(self.process_chunk, chunks)
# 合并结果
if isinstance(data, (list, tuple)):
final_result = []
for chunk_result in results:
final_result.extend(chunk_result)
else:
final_result = pd.concat(results, ignore_index=True)
return final_result
# 使用示例
def parallel_data_processing():
# 生成测试数据
data_size = 10000
data = list(range(data_size))
# 创建处理器
processor = BatchDataProcessor(chunk_size=2000, num_workers=4)
# 执行批量处理
result = processor.batch_process(data)
print(f"原始数据大小: {len(data)}")
print(f"处理后数据大小: {len(result)}")
print(f"前10个结果: {result[:10]}")
return result
# 运行示例
# processed_data = parallel_data_processing()
带状态监控的进程池
from multiprocessing import Pool, Manager
import time
class ProcessPoolMonitor:
def __init__(self, num_workers=4):
self.num_workers = num_workers
self.manager = Manager()
self.status_dict = self.manager.dict()
self.progress_queue = self.manager.Queue()
def monitored_worker(self, task_id, func, *args, **kwargs):
"""带监控的工作函数"""
worker_name = f"Worker-{task_id}"
# 更新状态
self.status_dict[worker_name] = {
'task': task_id,
'status': 'running',
'progress': 0,
'start_time': time.time()
}
try:
# 执行任务(模拟进度更新)
result = func(*args, **kwargs)
# 更新完成状态
self.status_dict[worker_name] = {
'task': task_id,
'status': 'completed',
'progress': 100,
'end_time': time.time()
}
# 发送进度信息
self.progress_queue.put({
'worker': worker_name,
'task': task_id,
'status': 'completed'
})
return result
except Exception as e:
self.status_dict[worker_name] = {
'task': task_id,
'status': 'failed',
'error': str(e)
}
raise
def run_tasks(self, tasks):
"""运行任务列表"""
with Pool(self.num_workers) as pool:
# 异步提交所有任务
async_results = []
for i, (func, args, kwargs) in enumerate(tasks):
result = pool.apply_async(
self.monitored_worker,
args=(i, func) + args,
kwds=kwargs
)
async_results.append(result)
# 监控进度(非阻塞)
while True:
# 检查是否所有任务完成
all_done = all(result.ready() for result in async_results)
# 获取进度更新
try:
while True:
progress = self.progress_queue.get_nowait()
print(f"进度更新: {progress}")
except:
pass
# 显示当前状态
print("\n当前状态:")
for worker, status in self.status_dict.items():
print(f" {worker}: {status['status']}")
if all_done:
break
time.sleep(0.5)
# 收集结果
results = [result.get() for result in async_results]
return results
# 使用示例
def sample_task(name, duration):
"""示例任务"""
time.sleep(duration)
return f"任务 {name} 完成"
tasks = [
(sample_task, ('任务1', 2), {}),
(sample_task, ('任务2', 1.5), {}),
(sample_task, ('任务3', 3), {}),
(sample_task, ('任务4', 1), {}),
]
monitor = ProcessPoolMonitor(num_workers=2)
# results = monitor.run_tasks(tasks)
异常处理和重试机制
from multiprocessing import Pool
import random
class ResilientProcessPool:
def __init__(self, num_workers=4, max_retries=3):
self.num_workers = num_workers
self.max_retries = max_retries
def execute_with_retry(self, func, args, kwargs):
"""带重试的任务执行"""
for attempt in range(self.max_retries):
try:
return func(*args, **kwargs)
except Exception as e:
if attempt < self.max_retries - 1:
wait_time = 2 ** attempt # 指数退避
print(f"重试 {attempt + 1}/{self.max_retries}, 等待 {wait_time}秒")
time.sleep(wait_time)
else:
raise e
def safe_map(self, func, iterable):
"""安全批量执行"""
with Pool(self.num_workers) as pool:
# 创建任务列表
tasks = [(func, (item,), {}) for item in iterable]
# 提交任务
async_results = []
for task_func, args, kwargs in tasks:
result = pool.apply_async(
self.execute_with_retry,
args=(task_func, args, kwargs)
)
async_results.append(result)
# 收集结果,处理异常
results = []
for i, async_result in enumerate(async_results):
try:
result = async_result.get(timeout=10)
results.append(result)
except Exception as e:
print(f"任务 {i} 最终失败: {e}")
results.append(None)
return results
# 使用示例
def potentially_failing_task(x):
"""可能失败的任务"""
if random.random() < 0.3: # 30%概率失败
raise ValueError(f"随机失败: {x}")
return x * 2
resilient_pool = ResilientProcessPool(num_workers=3, max_retries=2)
# results = resilient_pool.safe_map(potentially_failing_task, range(20))
-
选择合适的进程数:
import os num_workers = os.cpu_count() # 通常使用CPU核心数
-
内存管理:
# 使用迭代器避免内存溢出 with Pool(4) as pool: results = pool.imap(process_func, large_dataset) for result in results: process_single_result(result) -
超时控制:
# 设置超时 try: result = async_result.get(timeout=30) except multiprocessing.TimeoutError: print("任务超时")
这些案例涵盖了Python进程池的主要使用场景,你可以根据实际需求选择合适的方案。