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

如何在多进程中启动ThreadPoolExecutor提升I/O任务性能?

I/O密集型任务的多进程+多线程优化实现

问题背景

处理I/O密集型任务时,遇到以下瓶颈:

  • 单进程启动1000个线程的ThreadPoolExecutor耗时极长,CLI会冻结2-3分钟才开始执行任务;
  • 单纯使用multiprocessing.Pool受限于8核CPU,只能同时运行8个任务,无法满足大量并发需求。

需要将线程分摊到多核进程中,结合多进程的CPU并行性和多线程的I/O并发优势,充分利用机器性能。

现有代码(单线程池实现)

from concurrent.futures import ThreadPoolExecutor, as_completed

class ThreadingxMultiprocessing():
    
    def __init__(self) -> None:
        
        profileTasks = ["TEST1",
                        "TEST2",
                        "TEST3",
                        "TEST4",
                        "TEST5",
                        "TEST6",
                        "TEST7",
                        "TEST8",
                        "TEST9",
                        "TEST10",
                        "TEST11",
                        "TEST12",
                        "TEST13",
                        "TEST14",
                        "TEST15",
                        "TEST16",
                        "TEST17",
                        "TEST18",
                        "TEST19",
                        "TEST20",
                        "TEST21",
                        "TEST22",
                        "TEST23",
                        "TEST24",
                        "... and some more to get to 1k profiles",]
        
        self.threads=1000
        
        while True:
                        
            with ThreadPoolExecutor(max_workers=self.threads) as executor:
                for index, profile in enumerate(profileTasks):
                    
                    executor.submit(
                        self.runTask, index, profile
                    )

            
            break
     
    def runTask(self, index, profile): 
        print(index,profile)

ThreadingxMultiprocessing()

思路代码(未完成)

用户已明确按CPU核心数拆分线程的思路,但缺少具体实现:

from concurrent.futures import ThreadPoolExecutor, as_completed
import multiprocessing
import math
number_of_cpucores = multiprocessing.cpu_count()

class ThreadingxMultiprocessing():
    
    def __init__(self) -> None:
        
        profileTasks = ["TEST1",
                        "TEST2",
                        "TEST3",
                        "TEST4",
                        "TEST5",
                        "TEST6",
                        "TEST7",
                        "TEST8",
                        "TEST9",
                        "TEST10",
                        "TEST11",
                        "TEST12",
                        "TEST13",
                        "TEST14",
                        "TEST15",
                        "TEST16",
                        "TEST17",
                        "TEST18",
                        "TEST19",
                        "TEST20",
                        "TEST21",
                        "TEST22",
                        "TEST23",
                        "TEST24",
                        "... and some more to get to 1k profiles"]
        
        self.threads=1000
        #round them to get an integer datatype
        threads_in_each_process = math.ceil(float(self.threads)/ float(number_of_cpucores))
        
        #-> and then starting the thread pools e.g. with 125 threads each if you have 8 cores
        multiprocessing.Process()
        

    def runTask(self, index, profile): 
        print(index,profile)

ThreadingxMultiprocessing()

完整实现方案

以下代码实现了多进程+多线程的混合模型,每个进程启动对应数量的线程处理任务子集:

from concurrent.futures import ThreadPoolExecutor
import multiprocessing
import math

def runTask(index, profile):
    """实际I/O任务处理逻辑,保持原有功能"""
    print(index, profile)

def process_task_batch(task_batch):
    """单个进程的任务处理函数:启动线程池处理分配到的任务批次"""
    # 线程数设为当前批次的任务数量,避免资源浪费
    threads_in_process = len(task_batch)
    with ThreadPoolExecutor(max_workers=threads_in_process) as executor:
        # 提交所有线程任务
        futures = [executor.submit(runTask, idx, profile) for idx, profile in task_batch]
        # 等待所有线程任务完成(可选,根据需求决定是否需要同步)
        for future in futures:
            future.result()

class ThreadingxMultiprocessing():
    def __init__(self) -> None:
        # 模拟1000个任务
        profileTasks = [f"TEST{i+1}" for i in range(1000)]
        total_tasks = len(profileTasks)
        cpu_core_count = multiprocessing.cpu_count()
        
        # 拆分任务:按CPU核心数均分任务,每个进程处理一个批次
        tasks_per_process = math.ceil(total_tasks / cpu_core_count)
        task_batches = []
        for start_idx in range(0, total_tasks, tasks_per_process):
            # 保留任务的原始索引和profile,确保处理逻辑与原代码一致
            end_idx = min(start_idx + tasks_per_process, total_tasks)
            batch = list(enumerate(profileTasks[start_idx:end_idx], start=start_idx))
            task_batches.append(batch)
        
        # 启动进程池,每个进程处理一个任务批次
        with multiprocessing.Pool(processes=cpu_core_count) as pool:
            pool.map(process_task_batch, task_batches)

if __name__ == "__main__":
    ThreadingxMultiprocessing()

关键实现说明

  1. 任务拆分逻辑:将1000个任务按CPU核心数均分,每个进程拿到一个独立的任务批次(如8核时每个批次125个任务),避免单进程启动大量线程的开销。
  2. 进程内线程池:每个进程启动对应批次大小的线程池,专注处理自己的任务子集,最大化I/O并发效率。
  3. 进程池管理:使用multiprocessing.Pool自动管理进程生命周期,无需手动创建和维护Process对象,简化代码逻辑。
  4. 序列化兼容:将任务处理函数runTask和进程入口函数process_task_batch定义在模块级别,确保能被多进程的序列化机制(pickle)正确处理,避免类方法序列化失败的问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 05:40:31