Python 3.9 Watchdog 触发文件事件时无法启动多进程处理的问题排查
Python watchdog文件监控多进程不生效问题排查
需求背景
使用Python的watchdog模块开发文件处理类,期望每次新文件事件触发时都启动独立进程处理文件,满足多源文件同时涌入时的并发处理需求。实际运行时syslog仅显示单条python3.9进程日志,所有日志PID完全一致,未实现多进程效果。
原有实现代码
fileMonitor.py 内容
import time from watchdog.observers.polling import PollingObserver from watchdog.events import RegexMatchingEventHandler import converter ## ## \brief Class to handle monitored events ## class LogFileEventHandler(RegexMatchingEventHandler): MONITOR_REGEX = [r'.*\.(gz|txt)$'] # 仅监控后缀为.gz或.txt的文件 IGNORE_REGEX = [r'.*/archive/*'] # 忽略*/archive/*路径下的所有事件 ### ### Public methods ### def __init__(self): super().__init__( regexes=self.MONITOR_REGEX, ignore_regexes=self.IGNORE_REGEX, ignore_directories=True, case_sensitive=False) self.cm = converter.ConverterManager() def on_created(self, event): self.cm.convertMemory(event.src_path) def on_moved(self, event): self.cm.convertMemory(event.dest_path) ## ## \brief filesystem change monitor class ## \note 因为监控的是网络文件系统,必须使用PollingObserver ## 没有操作系统原生API支持网络文件系统的变更通知 ## class LogFileMonitor: ### ### Public methods ### def __init__(self, monitorPath): self.monitorPath = monitorPath # 待监控的路径 self.handler = LogFileEventHandler() # 事件处理实例 self.observer = PollingObserver() # 监控实现方式 def run(self): self._start() # 初始化observer try: while True: time.sleep(1) # 进程挂起 except KeyboardInterrupt: self._stop() # 终止observer ### ### Private methods ### def _start(self): self.observer.schedule( # 配置observer event_handler=self.handler, path=self.monitorPath, recursive=True, ) self.observer.start() # 启动observer def _stop(self): self.observer.stop() # 停止observer self.observer.join() # 等待observer完全退出
converter.py 原有内容
import time import concurrent.futures ## ## \brief 日志数据内存解析、写入influxdb的实现类 ## class ConverterMemoryWorker: ### ### Public methods ### def __init__(self, logFile): self.logFile = logFile def run(self): time.sleep(30) # 模拟长耗时处理逻辑 ## ## \brief 转换工作进程管理类 ## class ConverterManager: ### ### Public methods ### def __init__(self): print('Created new instance of ConverterManager') def convertMemory(self, logFile): with concurrent.futures.ProcessPoolExecutor(max_workers=4) as executor: # 从进程池创建新进程 executor.submit(self._task(logFile)) # 启动工作进程 ### ### Private methods ### def _task(self, logFile): converterWorker = ConverterMemoryWorker(logFile) converterWorker.run()
核心问题点
ProcessPoolExecutor创建逻辑错误:原有代码在convertMemory方法中每次调用都通过with语句创建新的进程池,with上下文会在退出前自动阻塞等待所有任务执行完成,等于每次处理事件都要等当前任务跑完才能处理下一个,完全无法并发。submit方法传参错误:原有代码写为executor.submit(self._task(logFile)),会先在主进程同步执行self._task(logFile)方法,再把执行结果提交给进程池,所有业务逻辑都在主进程运行,自然不会产生新的子进程。- 实例方法序列化问题:未做特殊处理的实例方法直接提交给进程池时,需要序列化整个类实例,容易出现pickle序列化失败的问题。
修复方案
- 将
ProcessPoolExecutor改为类初始化时全局创建,不要每次处理事件都新建进程池,复用进程资源。 - 修正
submit传参方式,改为executor.submit(可调用对象, 参数1, 参数2...)的格式,不要提前执行方法。 - 将任务执行方法
_task改为静态方法,避免序列化整个类实例的问题。
修复后可运行代码
import time import concurrent.futures ## ## \brief 日志数据内存解析、写入influxdb的实现类 ## class ConverterMemoryWorker: ### ### Public methods ### def __init__(self, logFile): self.logFile = logFile def run(self): print(f'Started process for {self.logFile}') time.sleep(10) # 模拟长耗时处理逻辑 print(f'Terminated process for {self.logFile}') ## ## \brief 转换工作进程管理类 ## class ConverterManager: executor = None ### ### Public methods ### def __init__(self): print('Created new instance of ConverterManager') self.executor = concurrent.futures.ProcessPoolExecutor(max_workers=4) def convertMemory(self, logFile): self.executor.submit(self._task, logFile) # 启动工作进程 ### ### Private methods ### @staticmethod def _task(logFile): converterWorker = ConverterMemoryWorker(logFile) converterWorker.run() if __name__ == '__main__': cm = ConverterManager() for i in range(30): cm.convertMemory(i)
内容的提问来源于stack exchange,提问作者Benny
相关产品推荐
相关产品推荐

