Python进程池案例如何批量管理进程

wen python案例 28

本文目录导读:

Python进程池案例如何批量管理进程

  1. 基础进程池使用
  2. 使用回调函数管理进度
  3. 动态任务队列管理
  4. 批量数据处理框架
  5. 带状态监控的进程池
  6. 异常处理和重试机制

我来介绍几种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))
  1. 选择合适的进程数

    import os
    num_workers = os.cpu_count()  # 通常使用CPU核心数
  2. 内存管理

    # 使用迭代器避免内存溢出
    with Pool(4) as pool:
        results = pool.imap(process_func, large_dataset)
        for result in results:
            process_single_result(result)
  3. 超时控制

    # 设置超时
    try:
        result = async_result.get(timeout=30)
    except multiprocessing.TimeoutError:
        print("任务超时")

这些案例涵盖了Python进程池的主要使用场景,你可以根据实际需求选择合适的方案。

抱歉,评论功能暂时关闭!