如何基于资源使用情况追踪并终止Python进程池中的特定进程
宏基因组pipeline进程池动态缩容实现方案
方案1:基于原生multiprocessing.Pool实现
原生Pool未提供公开的动态缩容接口,杀子进程后自动补进程是因为进程池的目标进程数属性仍为初始值5,调整逻辑如下:
- 进程池创建后,获取内部维护的子进程列表,按启动时间排序
- 触发缩容条件时,先修改进程池的私有属性
_processes为目标值3,阻止进程池补充新的子进程 - 终止启动时间最晚的2个活跃子进程,剩余3个进程会继续执行剩余任务
代码示例:
import time import threading import glob from multiprocessing import Pool def run_prodigal(directory): # 原有业务逻辑 pass def check_over_usage(): # 原有服务器wa值检查逻辑,返回True表示超过阈值 pass def multi_prodigal_processing(root): directories = glob.glob(root) prodigal_pool = Pool(processes=5) # 按启动时间排序子进程,列表默认顺序即为启动顺序 worker_processes = sorted(prodigal_pool._pool, key=lambda p: p._create_time) # 启动后台监控线程 threading.Thread( target=monitor_prodigal_usage, args=(prodigal_pool, worker_processes), daemon=True ).start() for directory in directories: prodigal_pool.apply_async(run_prodigal, (directory,)) prodigal_pool.close() prodigal_pool.join() def monitor_prodigal_usage(pool, worker_processes): scaled_down = False while True: if not scaled_down and check_over_usage(): # 修改进程池目标进程数,避免杀完自动补 pool._processes = 3 # 终止最晚启动的2个进程 for p in worker_processes[-2:]: if p.is_alive(): p.terminate() scaled_down = True # 所有任务执行完成后退出监控 if pool._state == Pool.CLOSE: break time.sleep(1)
注意:
_pool、_processes、_create_time为Pool私有属性,Python 3.8+版本测试可用,低版本需要对应调整属性名。
方案2:自定义Process+信号量实现(生产环境推荐)
避免依赖Pool的私有属性(不同Python版本私有属性命名可能变动),自行实现进程管控逻辑,稳定性更高:
- 用
multiprocessing.Semaphore控制并发上限,初始值设为5 - 用列表存储所有进程对象和对应的启动时间戳,方便后续按启动时间筛选要终止的进程
- 触发缩容条件时,修改信号量的内部值为3,再终止最晚启动的2个活跃进程即可,不会自动补充新进程
代码示例:
import time import threading import glob from multiprocessing import Process, Semaphore # 初始并发上限为5 max_concurrency = Semaphore(5) # 存储(进程对象,启动时间戳) process_records = [] def run_prodigal(directory): with max_concurrency: # 原有业务逻辑 pass def check_over_usage(): # 原有服务器wa值检查逻辑 pass def monitor_prodigal_usage(): scaled_down = False while True: if not scaled_down and check_over_usage(): # 调整并发上限为3 max_concurrency._value = 3 # 按启动时间排序,终止最晚启动的2个活跃进程 process_records.sort(key=lambda x: x[1]) for p, _ in process_records[-2:]: if p.is_alive(): p.terminate() scaled_down = True # 所有进程执行完成后退出监控 if all(not p.is_alive() for p, _ in process_records): break time.sleep(1) def multi_prodigal_processing(root): directories = glob.glob(root) # 启动后台监控线程 threading.Thread(target=monitor_prodigal_usage, daemon=True).start() # 启动所有进程 for directory in directories: p = Process(target=run_prodigal, args=(directory,)) p.start() process_records.append((p, time.time())) # 等待所有进程执行完成 for p, _ in process_records: p.join()
轻量替代方案
如果不需要严格终止已启动的进程,仅需限制后续并发数,可改用concurrent.futures.ProcessPoolExecutor,触发缩容时直接取消所有未启动的pending任务,同时将后续提交任务的并发阈值限制为3即可,实现成本更低。
内容的提问来源于stack exchange,提问作者Ron Turetzky
相关产品推荐
相关产品推荐

