Python多进程:如何限制Pool的内存使用量?
为Python进程池设置内存阈值控制进程启动的方案
Python标准库的multiprocessing.Pool本身没有内置的内存阈值配置项,但可以通过自定义监控逻辑+第三方工具实现类似效果,核心思路是在启动新进程前检查内存使用情况,满足阈值条件再提交任务。以下是两种实用方案:
方案1:监控子进程总内存,达标才启动新任务
通过psutil库获取所有子进程的内存占用总和,每次提交任务前检查该总和是否低于设定阈值,不足则等待重试。
步骤1:安装依赖
pip install psutil
步骤2:自定义内存控制进程池
import multiprocessing import psutil import time from typing import Callable class MemoryControlledPool: def __init__(self, max_workers: int, memory_threshold_mb: int): self.max_workers = max_workers # 转换为字节单位方便计算 self.memory_threshold = memory_threshold_mb * 1024 * 1024 self.pool = multiprocessing.Pool(max_workers=max_workers) # 记录活跃子进程 self.active_processes = [] def _calc_total_child_memory(self) -> int: """计算所有活跃子进程的内存占用总和""" total_rss = 0 to_remove = [] for proc in self.active_processes: try: # 获取进程的常驻内存大小(RSS) total_rss += proc.memory_info().rss except psutil.NoSuchProcess: # 进程已结束,标记移除 to_remove.append(proc) # 清理已结束的进程记录 for proc in to_remove: self.active_processes.remove(proc) return total_rss def apply_async(self, func: Callable, args: tuple = (), callback: Callable = None): """带内存检查的异步任务提交""" while True: current_total = self._calc_total_child_memory() if current_total < self.memory_threshold: # 内存充足,提交任务并记录子进程 def wrapped_task(*task_args): # 注册当前子进程到监控列表 self.active_processes.append(psutil.Process()) return func(*task_args) self.pool.apply_async(wrapped_task, args, callback) break # 内存不足,等待0.5秒后重试 time.sleep(0.5) def close(self): self.pool.close() self.pool.join() # 使用示例 def memory_hungry_task(task_id): # 模拟内存密集型操作:创建大列表 large_data = [0] * 10**7 time.sleep(2) return f"任务{task_id}完成,占用内存约{psutil.Process().memory_info().rss//(1024*1024)}MB" if __name__ == "__main__": # 配置:最大4个工作进程,子进程总内存阈值2048MB(2GB) pool = MemoryControlledPool(max_workers=4, memory_threshold_mb=2048) for i in range(10): pool.apply_async(memory_hungry_task, args=(i,), callback=lambda res: print(res)) pool.close()
方案2:监控系统可用内存,动态控制并发
如果不需要精确统计子进程内存,也可以监控系统剩余可用内存,当剩余内存低于阈值时,暂停启动新任务,适合对内存占用预估模糊的场景。
代码实现
import multiprocessing import psutil import time from threading import Semaphore def get_available_memory_mb() -> int: """获取系统当前可用内存(MB)""" mem_stats = psutil.virtual_memory() return mem_stats.available // (1024 * 1024) class DynamicConcurrencyPool: def __init__(self, max_workers: int, min_available_memory_mb: int): self.max_workers = max_workers self.min_free_mem = min_available_memory_mb # 用信号量控制最大并发数 self.semaphore = Semaphore(max_workers) self.pool = multiprocessing.Pool(max_workers=max_workers) def apply_async(self, func, args=(), callback=None): def wrapped_task(*task_args): with self.semaphore: # 等待系统可用内存达标 while get_available_memory_mb() < self.min_free_mem: time.sleep(0.5) return func(*task_args) return self.pool.apply_async(wrapped_task, args, callback) def close(self): self.pool.close() self.pool.join() # 使用示例 if __name__ == "__main__": # 配置:最大4个工作进程,系统可用内存至少保留1024MB(1GB)才启动新任务 pool = DynamicConcurrencyPool(max_workers=4, min_available_memory_mb=1024) for i in range(10): pool.apply_async(memory_hungry_task, args=(i,), callback=lambda res: print(res)) pool.close()
注意事项
- 两种方案都依赖
psutil获取内存数据,需要确保目标环境已安装该库; - 内存检查的间隔时间(
time.sleep(0.5))可根据实际场景调整,过短会增加CPU开销,过长会降低任务调度效率; - 方案1的子进程记录需要及时清理已结束的进程,避免内存统计错误;
- 如果任务内存波动大,建议预留足够的阈值缓冲空间,避免频繁触发等待逻辑。
内容的提问来源于stack exchange,提问作者ideasman42
相关产品推荐
相关产品推荐

