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

如何结合多线程与Watchdog实现文件夹新增文件的高效并行处理

解决方案

核心实现逻辑:使用Python标准库的concurrent.futures.ThreadPoolExecutor实现可控并发的文件处理,通过限制最大并行数避免内存超限,同时解决watchdog回调串行阻塞的问题。

修改后完整代码

import time
import threading
import cv2
import os
import csv
from concurrent.futures import ThreadPoolExecutor
from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler

folder_to_watch = "/trigger_folder"
# 最大并行数,可根据自身内存/CPU性能调整
MAX_WORKERS = 8
# CSV写入锁,避免多线程同时写文件导致内容错乱
csv_lock = threading.Lock()
# 你的CSV文件路径
CSV_PATH = "image_info.csv"

class EventHandler(FileSystemEventHandler):
    def __init__(self):
        # 初始化线程池
        self.executor = ThreadPoolExecutor(max_workers=MAX_WORKERS)
        # 初始化CSV表头(如果文件不存在的话)
        if not os.path.exists(CSV_PATH):
            with open(CSV_PATH, 'w', newline='', encoding='utf-8') as f:
                writer = csv.writer(f)
                writer.writerow(['文件路径', '宽度', '高度'])
    
    def do_smth(self, file_path):
        # 跳过文件夹创建事件
        if not os.path.isfile(file_path):
            return
        # 判断是否为图片文件,可根据需求补充后缀
        img_suffix = ('.jpg', '.jpeg', '.png', '.bmp')
        if not file_path.lower().endswith(img_suffix):
            return
        # 等待文件写入完成,避免触发事件时文件还没写完导致cv2读取失败
        while True:
            try:
                with open(file_path, 'rb') as f:
                    break
            except PermissionError:
                time.sleep(0.5)
        # 读取图片信息
        img = cv2.imread(file_path)
        if img is None:
            return
        height, width = img.shape[:2]
        # 写入CSV,加锁保证线程安全
        with csv_lock:
            with open(CSV_PATH, 'a', newline='', encoding='utf-8') as f:
                writer = csv.writer(f)
                writer.writerow([file_path, width, height])
        # 其他耗时操作写在这里
        print(f"处理完成:{file_path}")

    def on_created(self, event): # 当文件/文件夹创建时触发
        print(f"捕获到创建事件:{event.src_path}")
        # 提交任务到线程池,不会阻塞当前回调
        self.executor.submit(self.do_smth, event.src_path)

if __name__ == "__main__":
    observer = Observer()
    event_handler = EventHandler()
    observer.schedule(event_handler, path=folder_to_watch, recursive=False) # recursive设为True可监控子文件夹
    observer.start()

    try:
        while True:
            time.sleep(1)
    except KeyboardInterrupt:
        # 收到终止信号后先关闭线程池,等待所有已提交的任务处理完成再退出
        event_handler.executor.shutdown(wait=True)
        observer.stop()

    observer.join()

关键说明

  • 并发数控制:调整MAX_WORKERS参数即可控制同时处理的文件数量,数值越高处理速度越快,但内存占用也越高,建议根据单文件处理的内存占用和可用内存空间计算合适的数值。
  • 线程安全:涉及共享资源(比如示例中的CSV文件)的写操作必须加锁,避免多线程同时写入导致的内容错乱、文件损坏问题。
  • 异常兼容:代码中增加了文件写入等待逻辑,避免watchdog触发创建事件时文件还未完全写入磁盘,导致cv2读取失败的问题。

内容的提问来源于stack exchange,提问作者John_Doe

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 11:39:03