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

Python WatchDog库仅触发一次文件创建事件的问题求助

问题原因及解决方案

核心问题

你的WatchDog只能检测一次文件创建,根源在于on_created事件处理逻辑导致观察者线程阻塞或崩溃:

  1. data_analysis是耗时操作,直接在WatchDog的事件线程中执行会占用线程资源,导致后续事件无法被处理;
  2. 文件刚创建时可能还在写入过程中,此时读取会触发未捕获异常,直接终止观察者线程;
  3. 你把data_analysis嵌套在on_created内部,每次触发事件都会重新定义函数,额外增加不必要开销。

修复方案

1. 重构代码分离耗时操作

把data_analysis从on_created中移出,定义为独立函数,避免重复定义;

2. 多线程执行分析任务

用threading模块让分析任务在后台线程运行,不阻塞WatchDog的事件处理线程;

3. 添加异常捕获

防止单个文件处理失败导致整个监控进程终止;

4. 处理文件写入延迟

文件创建事件触发时,等待文件稳定后再读取,避免读取未完成的文件。

完整修复后的代码示例

from tkinter import *
from tkinter import filedialog
from watchdog.observers import Observer
from watchdog.events import PatternMatchingEventHandler
import pandas as pd
import numpy as np
import threading
import time

def data_analysis(src_path, logfunc=print):
    # 等待文件写入完成
    while True:
        try:
            with open(src_path, 'r') as f:
                break
        except:
            time.sleep(0.1)

    try:
        readdata = pd.read_csv(src_path, delimiter='\t', encoding="latin1", skiprows=24)
        df = pd.DataFrame(readdata)

        df = df.drop(labels=0, axis=0)
        
        df['Station']=df['Station'].astype(float).astype(int)

        # 初始化统计列
        stat_cols = [
            ("Axial Force", "Fz 1", 2600),
            ("Flexion", "FLPt", 58),
            ("IE", "IEPt", 5.7),
            ("AP", "APPt", 5.2)
        ]
        for name, _, _ in stat_cols:
            df[f"{name} Occurences"] = 0
            df[f"{name} Actual Value"] = pd.NaT

        # 数据类型转换
        float_cols = ['Fz 1','VLWf','FLPt','FLWf','IEPt','IEWf','APPt','APWf']
        for col in float_cols:
            df[col] = df[col].astype(float)
        int_cols = ['Fz 1','VLWf']
        for col in int_cols:
            df[col] = df[col].astype(int)

        # 筛选Station=1的数据
        data = df.loc[df['Station'] == 1]
        y = len(data.index)
        num = int(y * 0.03)

        # 首尾数据拼接
        first_rows = data.iloc[:num]
        last_rows = data.iloc[y-num:]
        data = pd.concat([last_rows, data, first_rows])
        z = len(data.index)
        data2 = data.set_index(np.linspace(1, z, z).astype(int))

        # 循环处理各统计项
        for name, col_name, tolerance_base in stat_cols:
            occur_list = []
            for i in range(num, z-num):
                val = data2[col_name].iloc[i]
                window = data2.iloc[i-num:i+num, float_cols.index(col_name)]
                lower = window - 0.05 * tolerance_base
                upper = window + 0.05 * tolerance_base
                
                if np.any(val >= lower) and np.any(val <= upper):
                    data2.at[i, f"{name} Occurences"] = 0
                else:
                    data2.at[i, f"{name} Occurences"] = 1
                    data2.at[i, f"{name} Actual Value"] = val
                    occur_list.append(i)
            
            total = data2[f"{name} Occurences"].sum()
            msg = f'The number of {name} values outside of the tolerance is: {total}'
            logfunc(msg)
            print(msg)
    except Exception as e:
        error_msg = f'Failed to process {src_path}: {str(e)}'
        logfunc(error_msg)
        print(error_msg)

class Watchdog(PatternMatchingEventHandler, Observer):
    def __init__(self, path='.', patterns='*', logfunc=print):
        PatternMatchingEventHandler.__init__(self, patterns)
        Observer.__init__(self)
        self.schedule(self, path=path, recursive=False)
        self.log = logfunc

    def on_created(self, event):
        if not event.is_directory:
            self.log(f"hey, {event.src_path} has been created!")
            # 启动后台线程执行分析
            threading.Thread(target=data_analysis, args=(event.src_path, self.log), daemon=True).start()

    def on_deleted(self, event):
        self.log(f"what the f**k! Someone deleted {event.src_path}!")

    def on_modified(self, event):
        self.log(f"hey buddy, {event.src_path} has been modified")

    def on_moved(self, event):
        self.log(f"ok ok ok, someone moved {event.src_path} to {event.dest_path}")

class GUI:
    def __init__(self):
        self.watchdog = None
        self.watch_path = '.'
        self.root = Tk()
        self.messagebox = Text(width=80, height=10)
        self.messagebox.pack()
        frm = Frame(self.root)
        Button(frm, text='Browse', command=self.select_path).pack(side=LEFT)
        Button(frm, text='Start Watchdog', command=self.start_watchdog).pack(side=RIGHT)
        Button(frm, text='Stop Watchdog', command=self.stop_watchdog).pack(side=RIGHT)
        frm.pack(fill=X, expand=1)
        self.root.mainloop()

    def start_watchdog(self):
        if self.watchdog is None:
            self.watchdog = Watchdog(path=self.watch_path, logfunc=self.log)
            self.watchdog.start()
            self.log('Watchdog started')
        else:
            self.log('Watchdog already started')

    def stop_watchdog(self):
        if self.watchdog:
            self.watchdog.stop()
            self.watchdog.join()
            self.watchdog = None
            self.log('Watchdog stopped')
        else:
            self.log('Watchdog is not running')

    def select_path(self):
        path = filedialog.askdirectory()
        if path:
            self.watch_path = path
            self.log(f'Selected path: {path}')

    def log(self, message):
        self.messagebox.insert(END, f'{message}\n')
        self.messagebox.see(END)      

if __name__ == '__main__':
    GUI()

关键修复点说明

  • 将data_analysis改为独立函数,消除重复定义的开销;
  • 用后台线程执行分析任务,避免阻塞WatchDog的事件处理线程;
  • 添加文件写入等待逻辑,确保读取完整文件;
  • 全局异常捕获,防止单个文件处理失败导致监控终止;
  • 过滤文件夹创建事件,只处理文件;
  • 停止WatchDog时调用join(),确保线程正确终止。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 00:00:59