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

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序列化失败的问题。

修复方案

  1. 将ProcessPoolExecutor改为类初始化时全局创建,不要每次处理事件都新建进程池,复用进程资源。
  2. 修正submit传参方式,改为executor.submit(可调用对象, 参数1, 参数2...)的格式,不要提前执行方法。
  3. 将任务执行方法_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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 10:54:01