You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

实现带工作进程数限制的多进程队列处理目录文件

实现带并发限制的目录文件处理队列

你可以用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 04:17:17