如何判断ProcessPoolExecutor是否已满?查询运行worker数量方法
查询ProcessPoolExecutor当前运行的Worker数量
下面提供几种实用方法来查看当前活跃的worker进程数量,判断是否达到max_workers上限:
方法1:直接访问私有属性(快速便捷)
ProcessPoolExecutor内部维护了_processes字典,存储当前活跃的worker进程,通过len(executor._processes)就能直接获取数量。
修改后的示例代码:
import concurrent.futures import time import threading def dummy_process(arg_a, arg_b): print("ml_process", arg_a, arg_b) time.sleep(5) executor = concurrent.futures.ProcessPoolExecutor(max_workers=2) def monitor_workers(): while True: active_count = len(executor._processes) print(f"当前活跃Worker数: {active_count}, 上限: {executor._max_workers}") time.sleep(1) def main(): # 启动后台监控线程 threading.Thread(target=monitor_workers, daemon=True).start() while True: executor.submit(dummy_process, "test_a", "test_b") time.sleep(0.5) # 控制提交速度,避免任务队列过载 if __name__ == "__main__": main()
注意:这是依赖Python内部私有API的方法,不同版本的Python可能会调整实现,适合快速调试场景。
方法2:自定义跨进程计数器(稳定可靠)
如果需要长期稳定的统计方式,可以用跨进程安全的计数器,在任务启动和结束时更新计数:
import concurrent.futures import time import multiprocessing import threading # 跨进程共享计数器,初始值为0 active_workers = multiprocessing.Value('i', 0) lock = multiprocessing.Lock() def dummy_process(arg_a, arg_b): with lock: active_workers.value += 1 try: print("ml_process", arg_a, arg_b) time.sleep(5) finally: with lock: active_workers.value -= 1 executor = concurrent.futures.ProcessPoolExecutor(max_workers=2) def monitor_workers(): while True: print(f"当前活跃Worker数: {active_workers.value}, 上限: {executor._max_workers}") time.sleep(1) def main(): threading.Thread(target=monitor_workers, daemon=True).start() while True: executor.submit(dummy_process, "test_a", "test_b") time.sleep(0.5) if __name__ == "__main__": main()
这种方法不依赖内部实现,兼容性更好,但需要修改任务函数的逻辑来维护计数。
方法3:通过Future对象追踪任务状态
可以维护一个集合保存未完成的Future对象,通过集合长度了解整体任务负载,但无法直接区分正在运行的worker和排队的任务:
import concurrent.futures import time import threading def dummy_process(arg_a, arg_b): print("ml_process", arg_a, arg_b) time.sleep(5) executor = concurrent.futures.ProcessPoolExecutor(max_workers=2) pending_futures = set() def monitor_workers(): while True: # 清理已完成的任务 done_futures = {f for f in pending_futures if f.done()} pending_futures.difference_update(done_futures) print(f"当前待处理/运行任务数: {len(pending_futures)}, Worker上限: {executor._max_workers}") time.sleep(1) def main(): threading.Thread(target=monitor_workers, daemon=True).start() while True: future = executor.submit(dummy_process, "test_a", "test_b") pending_futures.add(future) time.sleep(0.5) if __name__ == "__main__": main()
适合监控整体任务队列情况,若要精准获取运行中的worker数量,建议用前两种方法。
内容的提问来源于stack exchange,提问作者Ganessa
相关产品推荐
相关产品推荐

