You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python多进程任务仅获2倍提速,求8核CPU下优化方案

优化多进程资源利用的实用方法

你的8核CPU仅获得2倍提速,核心原因大概率是IO密集环节(读、写)与CPU密集环节(scipy处理)未解耦,导致CPU资源被IO等待浪费。以下是针对性优化方案:

1. 解耦IO与CPU任务,采用生产者-消费者模型

将读数据、处理数据、写数据拆分为独立阶段,用队列衔接不同任务,让IO操作和CPU计算并行执行,避免互相等待。

示例代码:

import concurrent.futures
from queue import Queue
import pandas as pd

# 假设read返回原始数据,treat接收原始数据并返回处理后的DataFrame
def read(time):
    # 实际读取逻辑
    return raw_data

def treat(time, raw_data):
    # 实际scipy处理逻辑
    return pd.DataFrame(processed_data)

def read_worker(schedule_queue, data_queue):
    while not schedule_queue.empty():
        time = schedule_queue.get()
        raw_data = read(time)
        data_queue.put((time, raw_data))
        schedule_queue.task_done()

def treat_worker(data_queue, result_queue):
    while True:
        time, raw_data = data_queue.get()
        if time is None:  # 任务结束信号
            data_queue.task_done()
            break
        df = treat(time, raw_data)
        result_queue.put((time, df))
        data_queue.task_done()

def write_worker(result_queue):
    while True:
        time, df = result_queue.get()
        if time is None:  # 任务结束信号
            result_queue.task_done()
            break
        df.to_csv(f'partial_{time}_data.csv')
        result_queue.task_done()

if __name__ == "__main__":
    schedule = [1,2,3,4,5,6,7,8]  # 示例时间列表
    schedule_queue = Queue()
    data_queue = Queue(maxsize=4)  # 限制队列大小,避免内存溢出
    result_queue = Queue(maxsize=4)

    # 填充任务队列
    for time in schedule:
        schedule_queue.put(time)

    # 启动IO读进程(2个,适配IO密集型任务)
    with concurrent.futures.ProcessPoolExecutor(max_workers=2) as read_exec:
        read_exec.submit(read_worker, schedule_queue, data_queue)
        
        # 启动CPU处理进程(8个,匹配核心数)
        with concurrent.futures.ProcessPoolExecutor(max_workers=8) as treat_exec:
            treat_workers = [treat_exec.submit(treat_worker, data_queue, result_queue) for _ in range(8)]
            
            # 启动IO写进程(2个)
            with concurrent.futures.ProcessPoolExecutor(max_workers=2) as write_exec:
                write_worker_fut = write_exec.submit(write_worker, result_queue)
                
                # 等待读任务完成
                schedule_queue.join()
                
                # 发送结束信号给处理进程
                for _ in range(8):
                    data_queue.put((None, None))
                data_queue.join()
                
                # 发送结束信号给写进程
                result_queue.put((None, None))
                result_queue.join()

2. 调整进程池大小,匹配任务类型

  • CPU密集型任务(treat):max_workers设置为CPU核心数(8)即可,过多进程会导致上下文切换开销激增。
  • IO密集型任务(read/write):可设置为核心数的2-4倍,但不要超过10,避免IO资源竞争。

你之前设置max_workers=16,对于8核CPU来说会加重上下文切换负担,建议先将process_map的max_workers改为8测试效果。

3. 优化CPU密集的treat函数

先提升单任务处理速度,再放大多进程收益:

  • 用numba对核心计算逻辑做JIT编译,加速数值运算。
  • 检查scipy函数是否有内置并行参数(如scipy.optimize的workers),开启内置并行。
  • 用numpy向量化操作替代Python循环,减少解释层开销。

示例(numba加速):

from numba import jit

@jit(nopython=True)
def core_calculation(data):
    # 这里是scipy处理中的核心数值计算逻辑
    result = ...
    return result

def treat(time, raw_data):
    processed = core_calculation(raw_data)
    # 其他scipy操作或转换为DataFrame
    df = pd.DataFrame(processed)
    return df

4. 优化IO操作效率

  • 读取数据时,优先用Parquet、Feather等二进制格式替代文本格式;若必须用CSV,使用pandas.read_csv(engine='pyarrow')加速读取。
  • 写入CSV时,用df.to_csv(engine='pyarrow'),比默认引擎快2-3倍;允许的话先写入Parquet,后续再批量转CSV,IO效率会大幅提升。
  • 批量处理IO:合并多个time的数据一次性读取,处理后再拆分写入,减少IO调用次数。

5. 避免不必要的进程间数据拷贝

若read读取的数据需要在进程间传递,使用multiprocessing.Array或multiprocessing.Manager的共享容器,避免完整数据拷贝。如果每个子进程独立处理自身time的数据(如当前thread函数逻辑),则无需额外处理,但要确保无隐式跨进程数据传递。

内容的提问来源于stack exchange,提问作者jAke

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 15:13:20