Python Watchdog自定义事件处理器触发异常问题排查
我之前在Windows环境下用watchdog处理批量CSV文件时,也碰到过一模一样的问题——没加业务逻辑时事件触发全正常,一加上pandas读文件的代码,就开始丢事件。结合你的代码和场景,我梳理了几个核心原因和对应的解决办法:
1. 事件处理线程阻塞导致队列溢出
Watchdog的Observer是基于线程工作的,默认情况下,事件处理(也就是你代码里的on_any_event方法)是在Observer的主线程中执行的。如果你的CSV处理逻辑比较耗时(比如读取大文件、做复杂的数据清洗),事件队列里的新事件会因为处理线程被占用而无法及时处理,当队列满了之后就会直接丢弃后续事件——这就是你60个文件只触发52次的核心原因之一。
解决办法:把耗时的业务逻辑放到单独的线程里执行,让事件处理线程快速响应新事件,别让它卡在IO操作上。比如用Python内置的threading.Thread封装你的CSV处理代码:
import threading # 保留你原有的其他导入代码 class Handler(FileSystemEventHandler): # 改用on_created替代on_any_event,更精准地监听文件创建事件 @staticmethod def on_created(event): if event.is_directory: return None pathlower = str(event.src_path).lower() if ".csv" in pathlower: print(f"Watchdog received created event - {event.src_path}") print("File is csv file") # 把CSV处理逻辑丢给单独线程,daemon=True保证线程随主程序退出 threading.Thread(target=process_csv, args=(event.src_path,), daemon=True).start() else: print("File is not csv file") # 单独抽离业务逻辑函数 def process_csv(file_path): try: # 这里放你原来的CSV处理代码 df = pd.read_csv(file_path, names=ALL_INI.columnlistbundle.split("|"), sep="|", encoding='latin-1', engine='python') arraySelectedColumnsBundle = ALL_INI.selectedcolumnsbundle.split(",") bundle_df = df[np.array(arraySelectedColumnsBundle)] # ... 其他后续处理逻辑 except Exception as e: # 异常处理,确保不会影响主线程 print(f"Failed to process {file_path}: {str(e)}") ErrorHandler(0, 'In CSV Process ', '-7', 'Exception ' + str(e), '', 1, '')
2. Windows文件系统的临时文件/重命名坑
在Windows里,粘贴或者生成文件时,系统经常会先创建一个临时文件(比如带.tmp后缀),等文件完全写入后再重命名成你要的CSV文件。这时候watchdog会捕获到临时文件的created事件,但不会自动捕获重命名后的事件——而你的代码只判断.csv后缀,自然就漏掉了真正的目标文件事件。
解决办法:
- 同时监听
renamed事件,处理重命名后的CSV文件; - 读取文件前加个小延迟或循环检查,确保文件已经完全写入磁盘(避免读取还在写入的文件导致PermissionError)。
修改后的Handler示例:
class Handler(FileSystemEventHandler): @staticmethod def on_created(event): # 保留之前的on_created逻辑,处理直接创建的CSV文件 if event.is_directory: return None pathlower = str(event.src_path).lower() if ".csv" in pathlower: print(f"Watchdog received created event - {event.src_path}") threading.Thread(target=process_csv, args=(event.src_path,), daemon=True).start() @staticmethod def on_renamed(event): # 处理重命名产生的CSV文件 if not event.is_directory: dest_path_lower = event.dest_path.lower() if ".csv" in dest_path_lower: print(f"Watchdog received renamed event - {event.src_path} -> {event.dest_path}") threading.Thread(target=process_csv, args=(event.dest_path,), daemon=True).start()
优化后的process_csv函数(确保文件可读取):
def process_csv(file_path): try: # 循环检查文件是否可以正常打开,避免读取未完成的文件 while True: try: with open(file_path, 'r', encoding='latin-1') as f: # 读一行测试是否能正常访问 f.readline() break except (PermissionError, IOError): # 文件还在写入,等待0.1秒再试 time.sleep(0.1) # 这里执行你的CSV处理逻辑 df = pd.read_csv(file_path, names=ALL_INI.columnlistbundle.split("|"), sep="|", encoding='latin-1', engine='python') arraySelectedColumnsBundle = ALL_INI.selectedcolumnsbundle.split(",") bundle_df = df[np.array(arraySelectedColumnsBundle)] except Exception as e: print(f"Failed to process {file_path}: {str(e)}") ErrorHandler(0, 'In CSV Process ', '-7', 'Exception ' + str(e), '', 1, '')
3. 异常处理的潜在问题
你的代码在on_any_event里加了try-except,但如果ErrorHandler本身抛出异常,可能会中断Observer线程的运行,导致后续事件完全不触发。另外,默认的异常处理没有打印详细的错误信息,不利于排查哪些事件处理失败了。
优化建议:
class Handler(FileSystemEventHandler): @staticmethod def on_created(event): if event.is_directory: return None try: pathlower = str(event.src_path).lower() if ".csv" in pathlower: print(f"Watchdog received created event - {event.src_path}") threading.Thread(target=process_csv, args=(event.src_path,), daemon=True).start() else: print(f"Ignored non-CSV file: {event.src_path}") except Exception as e: # 确保异常不会向上传递,同时打印详细信息 print(f"Error handling event {event.src_path}: {str(e)}") # 调用ErrorHandler,但要确保它不会抛出异常 try: ErrorHandler(0, 'In Observer ', '-7', 'Exception ' + str(e), '', 1, '') except Exception as eh_err: print(f"ErrorHandler failed: {str(eh_err)}")
4. 调整Observer的队列大小
默认情况下,Watchdog的Observer线程队列大小是有限的(和Python版本、操作系统有关),你可以在创建Observer时指定更大的队列大小,减少事件因为队列满而丢失的概率:
class OnMyWatch: watchDirectory = "." def __init__(self): # 指定更大的队列大小,比如1000,根据你的文件数量调整 self.observer1 = Observer(max_queue_size=1000) def run(self): self.observer1.schedule(Handler(), self.watchDirectory, recursive=True) self.observer1.start() try: while True: time.sleep(1) except KeyboardInterrupt: self.observer1.stop() print("Observer stopped by user") except Exception as e: self.observer1.stop() print(f"Observer stopped unexpectedly: {str(e)}") self.observer1.join()
总结
最关键的两个优化点是:
- 把耗时业务逻辑从事件处理线程剥离,用单独线程执行,避免阻塞事件队列;
- 处理Windows的重命名事件,避免漏掉临时文件转CSV的情况。
按上面的方法修改后,你应该能解决事件触发次数不足的问题了。
内容的提问来源于stack exchange,提问作者ApliDev

