如何在多进程中启动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()
关键实现说明
- 任务拆分逻辑:将1000个任务按CPU核心数均分,每个进程拿到一个独立的任务批次(如8核时每个批次125个任务),避免单进程启动大量线程的开销。
- 进程内线程池:每个进程启动对应批次大小的线程池,专注处理自己的任务子集,最大化I/O并发效率。
- 进程池管理:使用
multiprocessing.Pool自动管理进程生命周期,无需手动创建和维护Process对象,简化代码逻辑。 - 序列化兼容:将任务处理函数
runTask和进程入口函数process_task_batch定义在模块级别,确保能被多进程的序列化机制(pickle)正确处理,避免类方法序列化失败的问题。
内容的提问来源于stack exchange,提问作者realsuspection
相关产品推荐
相关产品推荐

