欢迎光临
我们一直在努力

Python 高效并行计算

Python 高效并行计算


📑 目录

  • 并行计算基础
  • Python 并行方案全景图
  • ProcessPoolExecutor 深度解析
  • 性能优化策略
  • 实战模式与最佳实践
  • 常见问题与陷阱
  • 进阶话题
  • 快速参考表

  • 1. 并行计算基础

    1.1 为什么需要并行计算?

    单核性能瓶颈
    • 摩尔定律放缓:单核CPU频率提升遇到物理极限
    • 多核成为标配:现代CPU普遍拥有4-64核心
    • 数据量爆炸式增长:科学计算、AI训练、大数据分析需求暴增
    实际收益

    # 示例:处理1000个热力学平衡计算任务
    # 单核顺序执行:1000 × 2秒 = 33分钟
    # 8核并行执行:1000 × 2秒 ÷ 8 = 4分钟(理论加速比 8×)


    1.2 并行 vs 并发 vs 分布式

    概念定义实现方式适用场景
    并发 (Concurrency) 多个任务交替执行,看起来同时进行 单核CPU时间片轮转 IO密集型(网络请求、文件读写)
    并行 (Parallelism) 多个任务真正同时执行 多核CPU同时运行 CPU密集型(数值计算、图像处理)
    分布式 (Distributed) 多台机器协同计算 网络通信 超大规模计算(大数据、深度学习)

    形象比喻:

    • 并发:一个厨师在多个锅之间来回切换(单核)
    • 并行:多个厨师同时做菜(多核)
    • 分布式:多个厨房同时运作(多机)

    1.3 Python 的 GIL 问题

    什么是 GIL?

    GIL (Global Interpreter Lock) 是 CPython 解释器的全局锁,同一时刻只允许一个线程执行 Python 字节码。

    GIL 的影响

    import threading
    import time

    def cpu_bound_task():
    """CPU密集型任务"""
    total = 0
    for i in range(10**7):
    total += i
    return total

    # ❌ 多线程无法加速CPU密集型任务(受GIL限制)
    start = time.time()
    threads = [threading.Thread(target=cpu_bound_task) for _ in range(4)]
    for t in threads: t.start()
    for t in threads: t.join()
    print(f"多线程耗时: {time.time() start:.2f}秒") # 约8秒(无加速)

    # ✅ 多进程可以绕过GIL
    from concurrent.futures import ProcessPoolExecutor
    start = time.time()
    with ProcessPoolExecutor(max_workers=4) as executor:
    list(executor.map(lambda x: cpu_bound_task(), range(4)))
    print(f"多进程耗时: {time.time() start:.2f}秒") # 约2秒(4×加速)

    结论
    • CPU密集型 → 使用多进程(ProcessPoolExecutor)
    • IO密集型 → 使用多线程(ThreadPoolExecutor)或协程(asyncio)

    2. Python 并行方案全景图

    2.1 标准库方案对比

    方案模块优势劣势适用场景
    多进程 multiprocessing 绕过GIL,真并行 内存开销大,启动慢 数值计算、科学模拟
    进程池 ProcessPoolExecutor 简洁API,自动管理 功能相对基础 批量任务处理
    多线程 threading 轻量,共享内存 受GIL限制 网络IO、文件IO
    线程池 ThreadPoolExecutor 简洁API 受GIL限制 爬虫、API调用
    协程 asyncio 高并发,低开销 生态不完善 高并发网络服务

    2.2 第三方库推荐

    库特点典型用法
    joblib 简化多进程,自动缓存 Parallel(n_jobs=4)(delayed(func)(i) for i in data)
    pathos 支持lambda和类方法 解决pickle限制
    ray 分布式计算框架 大规模机器学习
    dask 并行DataFrame 大数据分析

    3. ProcessPoolExecutor 深度解析

    3.1 基础架构

    from concurrent.futures import ProcessPoolExecutor

    # 创建进程池
    with ProcessPoolExecutor(max_workers=4) as executor:
    # 提交任务
    future = executor.submit(func, arg)
    # 或批量映射
    results = executor.map(func, iterable)

    内部机制:

    主进程

    ├─→ 创建进程池 (fork/spawn)
    │ ├─ Worker进程1 (等待任务)
    │ ├─ Worker进程2 (等待任务)
    │ └─ Worker进程3 (等待任务)

    ├─→ 任务队列 (Task Queue)
    │ [Task1, Task2, Task3, …]

    ├─→ 任务分发 (自动负载均衡)
    │ Worker1 ← Task1
    │ Worker2 ← Task2
    │ Worker3 ← Task3

    └─→ 结果收集 (Result Queue)
    [Result1, Result2, Result3, …]


    3.2 核心参数详解

    max_workers – 进程数量

    import os

    # 推荐配置
    num_workers = os.cpu_count() 1 # 保留1核给系统

    # 不同任务类型的建议
    # CPU密集型: cpu_count() 或 cpu_count() – 1
    # IO密集型: cpu_count() * 2~5(可超配)
    # 混合型: 实验确定最优值

    initializer – 初始化器(⭐重点)

    作用:每个worker进程启动时执行一次,用于加载共享资源。

    import numpy as np

    # 全局变量(worker进程间共享)
    global_data = None

    def init_worker(large_dataset):
    """初始化:加载大数据集"""
    global global_data
    global_data = large_dataset # 每个进程加载一次
    print(f"Worker {os.getpid()} 初始化完成")

    def process_task(index):
    """任务函数:直接使用全局数据"""
    return np.sum(global_data[index])

    # 使用初始化器
    data = np.random.rand(1000, 1000)
    with ProcessPoolExecutor(max_workers=4,
    initializer=init_worker,
    initargs=(data,)) as executor:
    results = list(executor.map(process_task, range(1000)))

    性能对比:

    # ❌ 不使用初始化器(每次任务都传递数据)
    # 1000次任务 × 1MB数据序列化 = 1GB开销 → 耗时10秒

    # ✅ 使用初始化器(数据加载4次,任务只传索引)
    # 4个进程 × 1MB数据加载 = 4MB开销 → 耗时1秒


    3.3 submit() vs map()

    submit() – 单任务提交

    # 适用场景:任务参数不同、需要异步处理
    futures = []
    with ProcessPoolExecutor(max_workers=4) as executor:
    # 提交不同类型的任务
    futures.append(executor.submit(func1, arg1))
    futures.append(executor.submit(func2, arg2))
    futures.append(executor.submit(func3, arg3))

    # 按需获取结果
    for future in futures:
    try:
    result = future.result(timeout=10) # 阻塞等待
    print(result)
    except Exception as e:
    print(f"任务失败: {e}")

    特点:

    • 返回 Future 对象
    • 可设置超时时间
    • 支持回调函数 add_done_callback()
    • 灵活但需要手动管理
    map() – 批量映射(⭐推荐)

    # 适用场景:批量处理相同类型的任务
    def square(x):
    return x ** 2

    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(square, range(100)))
    # 等价于: [square(0), square(1), …, square(99)]

    特点:

    • 自动保持输入输出顺序
    • 返回生成器(懒加载)
    • 代码简洁
    • 推荐用于批量任务
    对比总结
    特性submit()map()
    结果顺序 按完成顺序 按输入顺序
    异常处理 需手动检查 自动传播
    超时控制 支持 不支持(需转list)
    代码复杂度
    适用场景 异构任务 批量同构任务

    3.4 异常处理

    方式1:map() 自动传播异常

    def risky_task(x):
    if x == 5:
    raise ValueError("x 不能为 5")
    return x ** 2

    with ProcessPoolExecutor(max_workers=4) as executor:
    try:
    results = list(executor.map(risky_task, range(10)))
    except ValueError as e:
    print(f"捕获异常: {e}") # 会在 x=5 时抛出

    方式2:submit() 手动检查

    def safe_process(future):
    """安全处理Future结果"""
    try:
    return future.result()
    except Exception as e:
    return f"错误: {e}"

    with ProcessPoolExecutor(max_workers=4) as executor:
    futures = [executor.submit(risky_task, i) for i in range(10)]
    results = [safe_process(f) for f in futures]

    方式3:内部捕获(推荐)

    def robust_task(x):
    """任务内部处理异常"""
    try:
    if x == 5:
    raise ValueError("x 不能为 5")
    return x ** 2
    except Exception as e:
    return None # 或返回默认值/错误码

    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(robust_task, range(10)))
    valid_results = [r for r in results if r is not None]


    3.5 上下文管理与资源清理

    为什么使用 with 语句?

    # ❌ 不推荐:手动管理(容易忘记关闭)
    executor = ProcessPoolExecutor(max_workers=4)
    results = list(executor.map(func, data))
    executor.shutdown(wait=True) # 容易忘记!

    # ✅ 推荐:自动管理
    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(func, data))
    # 自动调用 executor.shutdown(wait=True)

    shutdown() 参数

    executor.shutdown(wait=True, cancel_futures=False)

    • wait=True:阻塞直到所有任务完成
    • wait=False:立即返回(后台继续执行)
    • cancel_futures=True:取消未开始的任务

    4. 性能优化策略

    4.1 任务粒度设计

    原则:避免过细或过粗

    import time

    def small_task(x):
    """粒度过小:0.001秒"""
    return x ** 2

    def large_task(x):
    """粒度过大:10秒"""
    time.sleep(10)
    return x ** 2

    # ❌ 过细任务:进程创建开销 > 计算时间
    # 10000个任务 × 0.001秒 = 10秒计算
    # 进程间通信开销 = 50秒 → 总耗时60秒(反而变慢!)

    # ❌ 过粗任务:负载不均衡
    # 4个任务 × 10秒,但最后一个任务完成时其他进程空闲

    # ✅ 合适粒度:0.1-1秒/任务
    def optimal_task(batch):
    """处理一批数据"""
    return [x ** 2 for x in batch]

    # 将10000个任务打包成100批
    batches = [range(i*100, (i+1)*100) for i in range(100)]
    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(optimal_task, batches))

    经验值
    • 任务耗时 < 0.01秒:考虑批处理
    • 任务耗时 0.1-10秒:理想区间
    • 任务耗时 > 60秒:考虑拆分或检查点机制

    4.2 数据序列化开销(pickle)

    问题:大对象传递开销巨大

    import numpy as np

    # ❌ 每次传递100MB数据
    def process_data(large_array):
    return np.sum(large_array)

    data = np.random.rand(10000, 1000) # 100MB
    with ProcessPoolExecutor(max_workers=4) as executor:
    # 1000次任务 × 100MB = 100GB序列化开销!
    results = list(executor.map(process_data, [data]*1000))

    解决方案1:使用初始化器

    global_data = None

    def init(data):
    global global_data
    global_data = data # 每个进程只加载一次

    def process_index(idx):
    return np.sum(global_data[idx])

    with ProcessPoolExecutor(max_workers=4,
    initializer=init,
    initargs=(data,)) as executor:
    # 只传索引,不传数据
    results = list(executor.map(process_index, range(1000)))

    解决方案2:共享内存(Python 3.8+)

    from multiprocessing import shared_memory

    # 创建共享内存
    shm = shared_memory.SharedMemory(create=True, size=data.nbytes)
    shared_array = np.ndarray(data.shape, dtype=data.dtype, buffer=shm.buf)
    shared_array[:] = data[:]

    def process_shared(shm_name, shape, dtype, idx):
    # 子进程连接共享内存
    shm = shared_memory.SharedMemory(name=shm_name)
    arr = np.ndarray(shape, dtype=dtype, buffer=shm.buf)
    return np.sum(arr[idx])

    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(
    lambda idx: process_shared(shm.name, data.shape, data.dtype, idx),
    range(1000)
    ))

    shm.close()
    shm.unlink()


    4.3 进程池大小选择

    CPU密集型任务

    import os

    # 基础公式
    optimal_workers = os.cpu_count()

    # 考虑超线程(Hyper-Threading)
    # Intel CPU: 物理核心 × 2 = 逻辑核心
    # 实际加速比通常只有 1.2-1.3×

    # 推荐配置
    num_workers = os.cpu_count() 1 # 保留1核给操作系统

    IO密集型任务

    # 可以超配(因为大部分时间在等待IO)
    num_workers = os.cpu_count() * 2 # 甚至 *5、*10

    # 示例:网络爬虫
    from concurrent.futures import ThreadPoolExecutor # IO密集用线程池
    with ThreadPoolExecutor(max_workers=50) as executor:
    results = executor.map(fetch_url, urls)

    混合型任务

    # 需要实验确定最优值
    import time

    def benchmark(num_workers):
    start = time.time()
    with ProcessPoolExecutor(max_workers=num_workers) as executor:
    list(executor.map(your_task, data))
    return time.time() start

    # 测试不同worker数量
    for n in [2, 4, 8, 16]:
    elapsed = benchmark(n)
    print(f"{n} workers: {elapsed:.2f}秒")


    4.4 内存管理

    问题:子进程内存累积

    # 子进程处理1000个任务,每个任务生成100MB临时数据
    # → 100GB内存占用(可能导致OOM)

    def memory_leak_task(x):
    large_temp = np.random.rand(10000, 1000) # 100MB
    result = np.sum(large_temp)
    # large_temp 不会被及时释放!
    return result

    解决方案1:显式删除

    def clean_task(x):
    large_temp = np.random.rand(10000, 1000)
    result = np.sum(large_temp)
    del large_temp # 显式删除
    return result

    解决方案2:定期重启worker

    from concurrent.futures import ProcessPoolExecutor

    # 每处理100个任务重启一次进程池
    batch_size = 100
    for i in range(0, len(data), batch_size):
    batch = data[i:i+batch_size]
    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(process_task, batch))
    # with块结束,进程池销毁,内存释放


    5. 实战模式与最佳实践

    5.1 CPU密集型任务:热力学平衡计算

    from concurrent.futures import ProcessPoolExecutor
    from tqdm import tqdm
    import numpy as np

    # 全局变量(进程间共享)
    global_database = None

    def init_worker(database_path):
    """初始化:加载热力学数据库"""
    global global_database
    # 假设这是一个耗时的加载过程
    global_database = load_thermodynamic_database(database_path)
    print(f"Worker {os.getpid()} 数据库加载完成")

    def calculate_equilibrium(composition):
    """计算单个成分的相平衡"""
    try:
    # 使用全局数据库,避免重复加载
    result = global_database.compute_equilibrium(composition)
    return result['phase_fraction']
    except Exception as e:
    return np.nan

    def batch_calculate(compositions, database_path, num_workers=None):
    """批量计算相平衡"""
    if num_workers is None:
    num_workers = os.cpu_count() 1

    print(f"启动 {num_workers} 个计算进程…")

    with ProcessPoolExecutor(max_workers=num_workers,
    initializer=init_worker,
    initargs=(database_path,)) as executor:
    # 使用tqdm显示进度
    results = list(tqdm(
    executor.map(calculate_equilibrium, compositions),
    total=len(compositions),
    desc="计算进度",
    unit="sample"
    ))

    return results

    # 使用示例
    compositions = [{'Fe': 0.7, 'Cr': 0.2, 'Ni': 0.1} for _ in range(1000)]
    results = batch_calculate(compositions, 'database.tdb')


    5.2 IO密集型任务:批量文件处理

    from concurrent.futures import ThreadPoolExecutor # IO密集用线程池
    import os
    from pathlib import Path

    def process_file(file_path):
    """处理单个文件"""
    try:
    with open(file_path, 'r') as f:
    content = f.read()
    # 数据处理逻辑
    processed = content.upper()
    # 写入结果
    output_path = file_path.replace('.txt', '_processed.txt')
    with open(output_path, 'w') as f:
    f.write(processed)
    return True
    except Exception as e:
    print(f"处理失败 {file_path}: {e}")
    return False

    def batch_process_files(directory, pattern='*.txt'):
    """批量处理文件"""
    files = list(Path(directory).glob(pattern))
    print(f"找到 {len(files)} 个文件")

    # IO密集型任务可以超配线程数
    with ThreadPoolExecutor(max_workers=20) as executor:
    results = list(tqdm(
    executor.map(process_file, files),
    total=len(files),
    desc="处理文件"
    ))

    success_count = sum(results)
    print(f"成功处理 {success_count}/{len(files)} 个文件")

    # 使用示例
    batch_process_files('/data/text_files')


    5.3 混合型任务:数据 ETL 流水线

    from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor
    import pandas as pd

    def extract(file_path):
    """提取:读取文件(IO密集)"""
    return pd.read_csv(file_path)

    def transform(df):
    """转换:数据清洗(CPU密集)"""
    # 复杂的数据处理逻辑
    df['new_column'] = df['value'].apply(lambda x: expensive_computation(x))
    return df

    def load(df, output_path):
    """加载:写入数据库(IO密集)"""
    df.to_sql('table', connection, if_exists='append')

    def etl_pipeline(file_paths):
    """混合型ETL流水线"""
    # 阶段1:并行提取(IO密集 → 线程池)
    with ThreadPoolExecutor(max_workers=10) as executor:
    raw_data = list(executor.map(extract, file_paths))

    # 阶段2:并行转换(CPU密集 → 进程池)
    with ProcessPoolExecutor(max_workers=4) as executor:
    transformed_data = list(executor.map(transform, raw_data))

    # 阶段3:并行加载(IO密集 → 线程池)
    with ThreadPoolExecutor(max_workers=10) as executor:
    list(executor.map(load, transformed_data,
    [f'output_{i}.db' for i in range(len(transformed_data))]))

    print("ETL流程完成")


    5.4 进度监控与日志

    使用 tqdm 监控进度

    from tqdm import tqdm

    # 方式1:包装 map()
    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(tqdm(
    executor.map(func, data),
    total=len(data),
    desc="处理中",
    unit="个",
    ncols=80 # 进度条宽度
    ))

    # 方式2:包装 submit()
    futures = []
    with ProcessPoolExecutor(max_workers=4) as executor:
    for item in data:
    futures.append(executor.submit(func, item))

    # 使用 tqdm 监控 future 完成情况
    for future in tqdm(futures, desc="等待结果"):
    result = future.result()

    自定义进度回调

    from concurrent.futures import as_completed

    def process_with_callback(item):
    result = expensive_function(item)
    return item, result

    def main():
    futures = {}
    with ProcessPoolExecutor(max_workers=4) as executor:
    for item in data:
    future = executor.submit(process_with_callback, item)
    futures[future] = item

    # 按完成顺序处理结果
    for future in tqdm(as_completed(futures), total=len(futures)):
    item = futures[future]
    try:
    original, result = future.result()
    print(f"✓ {original}{result}")
    except Exception as e:
    print(f"✗ {item} 失败: {e}")

    日志记录(多进程安全)

    import logging
    from logging.handlers import QueueHandler, QueueListener
    from multiprocessing import Queue

    def setup_logging():
    """配置多进程安全的日志"""
    log_queue = Queue()

    # 主进程监听器
    handler = logging.FileHandler('process.log')
    handler.setFormatter(logging.Formatter(
    '%(asctime)s – %(processName)s – %(levelname)s – %(message)s'
    ))
    listener = QueueListener(log_queue, handler)
    listener.start()

    return log_queue, listener

    def init_worker_with_logging(log_queue):
    """子进程初始化日志"""
    logger = logging.getLogger()
    logger.addHandler(QueueHandler(log_queue))
    logger.setLevel(logging.INFO)

    def logged_task(x):
    """带日志的任务"""
    logger = logging.getLogger()
    logger.info(f"开始处理 {x}")
    result = x ** 2
    logger.info(f"完成处理 {x}{result}")
    return result

    # 使用示例
    log_queue, listener = setup_logging()
    with ProcessPoolExecutor(max_workers=4,
    initializer=init_worker_with_logging,
    initargs=(log_queue,)) as executor:
    results = list(executor.map(logged_task, range(100)))
    listener.stop()


    6. 常见问题与陷阱

    6.1 Windows 下必须使用 if __name__ == '__main__'

    问题原因

    Windows 使用 spawn 方式创建子进程(导入主模块),而 Linux 使用 fork(复制进程)。

    # ❌ 错误示例(Windows会无限递归创建进程)
    from concurrent.futures import ProcessPoolExecutor

    def task(x):
    return x ** 2

    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(task, range(100)))
    # → Windows会报错: RuntimeError: An attempt has been made to start a new process…

    # ✅ 正确示例
    from concurrent.futures import ProcessPoolExecutor

    def task(x):
    return x ** 2

    if __name__ == '__main__': # 必须加这一行!
    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(task, range(100)))

    跨平台兼容写法

    import multiprocessing as mp

    if __name__ == '__main__':
    # 显式设置启动方式(推荐)
    mp.set_start_method('spawn') # 或 'fork', 'forkserver'

    # 其余代码…


    6.2 Lambda 函数不能被 pickle

    问题

    # ❌ Lambda无法序列化
    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(lambda x: x**2, range(100)))
    # → AttributeError: Can't pickle <lambda>

    解决方案

    # 方式1:使用命名函数
    def square(x):
    return x ** 2

    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(square, range(100)))

    # 方式2:使用 pathos 库(支持lambda)
    from pathos.multiprocessing import ProcessPool
    pool = ProcessPool(nodes=4)
    results = pool.map(lambda x: x**2, range(100))

    # 方式3:使用 partial
    from functools import partial

    def power(x, n):
    return x ** n

    with ProcessPoolExecutor(max_workers=4) as executor:
    square_func = partial(power, n=2)
    results = list(executor.map(square_func, range(100)))


    6.3 资源泄漏

    问题:忘记关闭进程池

    # ❌ 资源泄漏
    executor = ProcessPoolExecutor(max_workers=4)
    results = list(executor.map(task, data))
    # 忘记调用 executor.shutdown()
    # → 进程不会释放,内存泄漏

    解决方案

    # ✅ 使用 with 语句(推荐)
    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(task, data))
    # 自动关闭

    # ✅ 手动关闭
    executor = ProcessPoolExecutor(max_workers=4)
    try:
    results = list(executor.map(task, data))
    finally:
    executor.shutdown(wait=True)


    6.4 死锁与僵尸进程

    问题1:子进程等待主进程,主进程等待子进程

    from multiprocessing import Queue

    def task(queue):
    # 子进程等待主进程放入数据
    data = queue.get()
    return data * 2

    queue = Queue()
    with ProcessPoolExecutor(max_workers=1) as executor:
    future = executor.submit(task, queue)
    # 主进程等待子进程完成,但子进程在等待主进程
    result = future.result() # 死锁!
    queue.put(10)

    解决方案:先放数据

    queue = Queue()
    queue.put(10) # 先放数据
    with ProcessPoolExecutor(max_workers=1) as executor:
    future = executor.submit(task, queue)
    result = future.result()

    问题2:僵尸进程

    # 子进程已完成但主进程未回收
    executor = ProcessPoolExecutor(max_workers=4)
    futures = [executor.submit(task, i) for i in range(100)]
    # 未等待结果就退出 → 僵尸进程

    解决方案

    # 方式1:使用 wait()
    from concurrent.futures import wait
    wait(futures)

    # 方式2:获取所有结果
    results = [f.result() for f in futures]

    # 方式3:使用 with
    with ProcessPoolExecutor(max_workers=4) as executor:
    futures = [executor.submit(task, i) for i in range(100)]
    # with 块会自动等待所有任务完成


    6.5 共享状态问题

    问题:全局变量修改不共享

    counter = 0

    def increment():
    global counter
    counter += 1
    return counter

    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(lambda _: increment(), range(100)))

    print(counter) # 输出: 0(主进程的counter未改变)
    print(results) # 输出: [1, 1, 1, 1, …](每个子进程独立计数)

    解决方案:使用 Manager

    from multiprocessing import Manager

    def init_worker(shared_dict):
    global shared_data
    shared_data = shared_dict

    def increment_shared():
    with shared_data['lock']:
    shared_data['counter'] += 1
    return shared_data['counter']

    if __name__ == '__main__':
    manager = Manager()
    shared_dict = manager.dict({'counter': 0})
    shared_dict['lock'] = manager.Lock()

    with ProcessPoolExecutor(max_workers=4,
    initializer=init_worker,
    initargs=(shared_dict,)) as executor:
    results = list(executor.map(lambda _: increment_shared(), range(100)))

    print(shared_dict['counter']) # 输出: 100


    7. 进阶话题

    7.1 进程间通信(IPC)

    Queue – 线程/进程安全的队列

    from multiprocessing import Queue, Process

    def producer(queue):
    for i in range(10):
    queue.put(f"数据-{i}")
    queue.put(None) # 发送结束信号

    def consumer(queue):
    while True:
    item = queue.get()
    if item is None:
    break
    print(f"消费: {item}")

    if __name__ == '__main__':
    queue = Queue()
    p1 = Process(target=producer, args=(queue,))
    p2 = Process(target=consumer, args=(queue,))
    p1.start()
    p2.start()
    p1.join()
    p2.join()

    Pipe – 双向管道

    from multiprocessing import Pipe, Process

    def worker(conn):
    conn.send("Hello from worker")
    msg = conn.recv()
    print(f"Worker received: {msg}")
    conn.close()

    if __name__ == '__main__':
    parent_conn, child_conn = Pipe()
    p = Process(target=worker, args=(child_conn,))
    p.start()

    msg = parent_conn.recv()
    print(f"Parent received: {msg}")
    parent_conn.send("Hello from parent")

    p.join()


    7.2 共享内存(Python 3.8+)

    from multiprocessing import shared_memory
    import numpy as np

    # 创建共享内存
    shm = shared_memory.SharedMemory(create=True, size=1000000)
    arr = np.ndarray((100, 100), dtype=np.float64, buffer=shm.buf)
    arr[:] = np.random.rand(100, 100)

    def worker(shm_name, shape, dtype):
    """子进程访问共享内存"""
    shm = shared_memory.SharedMemory(name=shm_name)
    arr = np.ndarray(shape, dtype=dtype, buffer=shm.buf)
    result = np.sum(arr)
    shm.close()
    return result

    if __name__ == '__main__':
    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(
    lambda _: worker(shm.name, (100, 100), np.float64),
    range(4)
    ))

    print(results)
    shm.close()
    shm.unlink()

    优势:

    • 零拷贝(zero-copy)
    • 适合大规模NumPy数组
    • 性能远超pickle

    7.3 与 NumPy/Pandas 结合

    NumPy 数组并行处理

    import numpy as np
    from concurrent.futures import ProcessPoolExecutor

    def parallel_apply(func, arr, axis=0, num_workers=4):
    """沿指定轴并行应用函数"""
    splits = np.array_split(arr, num_workers, axis=axis)

    with ProcessPoolExecutor(max_workers=num_workers) as executor:
    results = list(executor.map(func, splits))

    return np.concatenate(results, axis=axis)

    # 使用示例
    data = np.random.rand(10000, 100)
    result = parallel_apply(lambda x: x ** 2, data, num_workers=4)

    Pandas DataFrame 并行处理

    import pandas as pd
    from concurrent.futures import ProcessPoolExecutor

    def parallel_apply_df(func, df, num_workers=4):
    """并行处理DataFrame"""
    splits = np.array_split(df, num_workers)

    with ProcessPoolExecutor(max_workers=num_workers) as executor:
    results = list(executor.map(func, splits))

    return pd.concat(results)

    # 使用示例
    df = pd.DataFrame({'A': range(10000), 'B': range(10000)})
    result = parallel_apply_df(lambda x: x['A'] ** 2, df, num_workers=4)


    7.4 分布式扩展:Ray

    # 安装: pip install ray

    import ray

    ray.init() # 初始化Ray集群

    @ray.remote
    def expensive_task(x):
    """Ray自动并行化的任务"""
    return x ** 2

    # 提交任务(返回Future对象)
    futures = [expensive_task.remote(i) for i in range(1000)]

    # 获取结果
    results = ray.get(futures)

    ray.shutdown()

    Ray 优势:

    • 自动分布式调度
    • 支持GPU加速
    • 适合大规模机器学习

    8. 快速参考表

    8.1 选择正确的并行方案

    # 决策树
    if 任务类型 == 'CPU密集':
    if 数据量 < 10GB:
    使用 ProcessPoolExecutor
    else:
    使用 Ray / Dask
    elif 任务类型 == 'IO密集':
    if 并发数 < 1000:
    使用 ThreadPoolExecutor
    else:
    使用 asyncio
    elif 任务类型 == '混合':
    分阶段使用不同方案

    8.2 API 速查表

    操作ProcessPoolExecutormultiprocessing.Pool
    创建 ProcessPoolExecutor(max_workers=4) Pool(processes=4)
    批量映射 executor.map(func, data) pool.map(func, data)
    单任务提交 executor.submit(func, arg) pool.apply_async(func, arg)
    获取结果 future.result() async_result.get()
    关闭 executor.shutdown() pool.close(); pool.join()

    8.3 性能优化 Checklist

    • 使用 initializer 预加载数据
    • 任务粒度 0.1-10秒/任务
    • 避免传递大对象(使用共享内存)
    • 进程数 = CPU核心数 – 1
    • 使用 with 自动管理资源
    • 异常处理在任务函数内部
    • 使用 tqdm 监控进度
    • 定期重启进程池(长时间运行)

    8.4 常见错误速查

    错误信息原因解决方案
    Can't pickle <lambda> Lambda无法序列化 使用命名函数或pathos
    RuntimeError: An attempt… Windows未加保护 添加 if __name__ == '__main__'
    BrokenProcessPool 子进程崩溃 检查任务函数内的异常
    TimeoutError 任务超时 增加 result(timeout=N)
    内存持续增长 内存泄漏 显式 del 或定期重启进程池

    9. 完整示例:生产级代码模板

    #!/usr/bin/env python3
    """
    高性能并行计算模板
    支持进度监控、异常处理、日志记录、资源管理
    """

    import os
    import logging
    from concurrent.futures import ProcessPoolExecutor, as_completed
    from tqdm import tqdm
    from typing import List, Callable, Any

    # 全局变量(worker进程共享)
    global_resource = None

    def init_worker(resource_path: str):
    """初始化worker进程"""
    global global_resource
    try:
    # 加载共享资源(数据库、模型等)
    global_resource = load_resource(resource_path)
    logging.info(f"Worker {os.getpid()} 初始化成功")
    except Exception as e:
    logging.error(f"Worker初始化失败: {e}")
    raise

    def process_task(task_data: Any) > Any:
    """处理单个任务"""
    try:
    # 使用全局资源
    result = global_resource.compute(task_data)
    return {'success': True, 'data': task_data, 'result': result}
    except Exception as e:
    logging.error(f"任务失败 {task_data}: {e}")
    return {'success': False, 'data': task_data, 'error': str(e)}

    def parallel_process(
    tasks: List[Any],
    resource_path: str,
    num_workers: int = None,
    show_progress: bool = True
    ) > List[Any]:
    """
    并行处理任务列表

    Args:
    tasks: 任务列表
    resource_path: 共享资源路径
    num_workers: worker数量(默认CPU核心数-1)
    show_progress: 是否显示进度条

    Returns:
    结果列表
    """
    if num_workers is None:
    num_workers = max(1, os.cpu_count() 1)

    logging.info(f"启动 {num_workers} 个worker处理 {len(tasks)} 个任务")

    results = []
    failed_count = 0

    with ProcessPoolExecutor(
    max_workers=num_workers,
    initializer=init_worker,
    initargs=(resource_path,)
    ) as executor:
    # 提交所有任务
    futures = {
    executor.submit(process_task, task): task
    for task in tasks
    }

    # 使用tqdm包装进度条
    iterator = as_completed(futures)
    if show_progress:
    iterator = tqdm(iterator, total=len(futures), desc="处理中")

    # 收集结果
    for future in iterator:
    task = futures[future]
    try:
    result = future.result(timeout=300) # 5分钟超时
    results.append(result)
    if not result['success']:
    failed_count += 1
    except Exception as e:
    logging.error(f"任务执行异常 {task}: {e}")
    results.append({
    'success': False,
    'data': task,
    'error': f'Execution error: {e}'
    })
    failed_count += 1

    logging.info(f"处理完成: 成功 {len(results)failed_count}, 失败 {failed_count}")
    return results

    def main():
    """主函数"""
    # 配置日志
    logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s – %(levelname)s – %(message)s',
    handlers=[
    logging.FileHandler('parallel_process.log'),
    logging.StreamHandler()
    ]
    )

    # 准备任务
    tasks = list(range(1000))
    resource_path = 'resource.db'

    # 执行并行处理
    results = parallel_process(
    tasks=tasks,
    resource_path=resource_path,
    num_workers=4,
    show_progress=True
    )

    # 保存结果
    success_results = [r for r in results if r['success']]
    print(f"成功处理: {len(success_results)}/{len(tasks)}")

    if __name__ == '__main__':
    main()


    10. 学习资源与延伸阅读

    官方文档

    • concurrent.futures 文档
    • multiprocessing 文档
    • PEP 3148 – futures 提案

    推荐书籍

    • 《Python并行编程实战》
    • 《高性能Python》

    第三方库

    • joblib – 简化并行计算
    • ray – 分布式计算框架
    • dask – 并行数据分析
    • pathos – 增强的并行工具

    性能分析工具

    • cProfile – 性能分析
    • line_profiler – 逐行分析
    • memory_profiler – 内存分析

    附录:完整代码索引

    A1. 基础示例

    # 最简单的并行计算
    from concurrent.futures import ProcessPoolExecutor

    def square(x):
    return x ** 2

    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(square, range(100)))

    A2. 带初始化器

    global_data = None

    def init(data):
    global global_data
    global_data = data

    def process(idx):
    return global_data[idx] ** 2

    with ProcessPoolExecutor(max_workers=4,
    initializer=init,
    initargs=(data,)) as executor:
    results = list(executor.map(process, range(len(data))))

    A3. 带进度条

    from tqdm import tqdm

    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(tqdm(
    executor.map(func, data),
    total=len(data),
    desc="处理中"
    ))

    A4. 异常处理

    def safe_task(x):
    try:
    return x ** 2
    except Exception as e:
    return None

    with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(safe_task, data))
    valid_results = [r for r in results if r is not None]


    赞(0)
    未经允许不得转载:171主机测评 » Python 高效并行计算
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址