实现带工作进程数限制的多进程队列处理目录文件
实现带并发限制的目录文件处理队列
你可以用multiprocessing.Queue实现任务队列,同时启动固定数量的工作进程持续从队列中获取任务执行,完全不需要修改原有的WorkOnFiles类。以下是完整的可运行代码:
import os import time import multiprocessing # 保留你原有的WorkOnFiles类,无需修改 class WorkOnFiles(multiprocessing.Process): def __init__(self, filename): super().__init__() self.filename = filename def run(self): # 这里写你原本的文件处理逻辑 print(f"开始处理文件: {self.filename}") time.sleep(3) # 模拟耗时操作 print(f"完成处理文件: {self.filename}") # 工作进程的任务函数:持续从队列取任务执行 def worker(queue): while True: filename = queue.get() # 收到None信号时退出进程 if filename is None: break # 实例化并运行原有的WorkOnFiles进程 worker_process = WorkOnFiles(filename) worker_process.run() if __name__ == "__main__": directory = "/path/to/directory" max_concurrent_workers = 2 # 限制同时处理的进程数 scan_interval = 1 # 目录扫描间隔(秒) # 创建进程安全的任务队列 task_queue = multiprocessing.Queue() # 启动指定数量的工作进程 workers = [] for _ in range(max_concurrent_workers): p = multiprocessing.Process(target=worker, args=(task_queue,)) p.start() workers.append(p) try: while True: # 扫描目录 for filename in os.listdir(directory): # 只处理ABC开头的文件 if filename.startswith("ABC"): old_path = os.path.join(directory, filename) new_filename = f"XYZ{filename}" new_path = os.path.join(directory, new_filename) # 用os.rename替代os.system,更安全且跨平台 os.rename(old_path, new_path) # 将绝对路径加入队列,避免子进程工作目录差异问题 task_queue.put(new_path) time.sleep(scan_interval) except KeyboardInterrupt: # 捕获Ctrl+C,发送终止信号给所有工作进程 print("\n正在终止所有工作进程...") for _ in range(max_concurrent_workers): task_queue.put(None) # 等待所有工作进程退出 for p in workers: p.join() print("所有进程已终止")
关键说明:
- 任务队列:使用
multiprocessing.Queue实现进程间安全的任务传递,主进程负责生产任务,工作进程负责消费任务。 - 并发限制:通过启动固定数量的工作进程(
max_concurrent_workers),确保同时处理的文件数不超过设定值。 - 原类兼容:完全保留
WorkOnFiles类的逻辑,只需在worker函数中实例化并调用其run方法。 - 安全重命名:用
os.rename替代os.system执行文件重命名,避免shell注入风险,同时支持跨平台。 - 优雅退出:捕获
KeyboardInterrupt信号,向队列发送None终止符,确保所有工作进程能安全退出。
内容的提问来源于stack exchange,提问作者jlipinski
相关产品推荐
相关产品推荐

