Available for day contractsFrom 21st September I have availability for day and half day contracts. Please contact for more information.

Contact →
mikepreston.org

Python Watchdog

File system monitoring library for tracking changes to files and directories in real-time.

Python Watchdog

File system monitoring library for tracking changes to files and directories in real-time.

Overview

Watchdog is a Python library that monitors file system events using platform-specific APIs (inotify on Linux, FSEvents on macOS, ReadDirectoryChangesW on Windows). It provides a clean API for observing directories and reacting to file creation, modification, deletion, and movement events.

WatchdogFile SystemEventsFiles/DirectoriesObserverEvent Handleron_createdon_modifiedon_deletedon_movedYour CodeWatchdogFile SystemEventsFiles/DirectoriesObserverEvent Handleron_createdon_modifiedon_deletedon_movedYour Code

Installation

pip install watchdog

Setting Up Observers and Event Handlers

Observers monitor directories whilst event handlers process detected changes.

File SystemEventHandlerSchedulerObserverFile SystemEventHandlerSchedulerObserverschedule(handler, path)start()File event detecteddispatch(event)on_modified(event)File SystemEventHandlerSchedulerObserverFile SystemEventHandlerSchedulerObserverschedule(handler, path)start()File event detecteddispatch(event)on_modified(event)

Basic Setup

from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler
import time

class MyHandler(FileSystemEventHandler):
    def on_any_event(self, event):
        print(f"{event.event_type}: {event.src_path}")

# Create and configure observer
observer = Observer()
handler = MyHandler()
observer.schedule(handler, path="/path/to/watch", recursive=True)

# Start monitoring
observer.start()
try:
    while True:
        time.sleep(1)
except KeyboardInterrupt:
    observer.stop()
observer.join()

Context Manager Pattern

from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler
from contextlib import contextmanager

@contextmanager
def watch_directory(path, handler, recursive=True):
    observer = Observer()
    observer.schedule(handler, path, recursive=recursive)
    observer.start()
    try:
        yield observer
    finally:
        observer.stop()
        observer.join()

# Usage
with watch_directory("/path/to/watch", MyHandler()) as observer:
    while True:
        time.sleep(1)

Multiple Directories

observer = Observer()

# Watch multiple paths with different handlers
observer.schedule(LogHandler(), "/var/log", recursive=False)
observer.schedule(ConfigHandler(), "/etc/app", recursive=True)
observer.schedule(DataHandler(), "/data", recursive=True)

observer.start()

Monitoring File System Events

Event Types

Event Type Description Attribute
FileCreatedEvent File created event.src_path
FileModifiedEvent File modified event.src_path
FileDeletedEvent File deleted event.src_path
FileMovedEvent File moved/renamed event.src_path, event.dest_path
DirCreatedEvent Directory created event.src_path
DirModifiedEvent Directory modified event.src_path
DirDeletedEvent Directory deleted event.src_path
DirMovedEvent Directory moved event.src_path, event.dest_path

Event Properties

def on_any_event(self, event):
    print(f"Event type: {event.event_type}")
    print(f"Is directory: {event.is_directory}")
    print(f"Source path: {event.src_path}")

    # For move events only
    if hasattr(event, 'dest_path'):
        print(f"Destination: {event.dest_path}")

Handling File Creation, Modification, and Deletion

Complete Event Handler

from watchdog.events import FileSystemEventHandler
import os

class ComprehensiveHandler(FileSystemEventHandler):
    def on_created(self, event):
        if event.is_directory:
            print(f"Directory created: {event.src_path}")
        else:
            print(f"File created: {event.src_path}")
            self.process_new_file(event.src_path)

    def on_modified(self, event):
        if not event.is_directory:
            print(f"File modified: {event.src_path}")
            self.process_modified_file(event.src_path)

    def on_deleted(self, event):
        if event.is_directory:
            print(f"Directory deleted: {event.src_path}")
        else:
            print(f"File deleted: {event.src_path}")
            self.cleanup_references(event.src_path)

    def on_moved(self, event):
        print(f"Moved: {event.src_path} -> {event.dest_path}")
        self.update_references(event.src_path, event.dest_path)

    def process_new_file(self, path):
        # Custom processing logic
        pass

    def process_modified_file(self, path):
        # Custom processing logic
        pass

    def cleanup_references(self, path):
        # Cleanup logic
        pass

    def update_references(self, old_path, new_path):
        # Update references logic
        pass

Debouncing Rapid Events

Files often trigger multiple modification events. Use debouncing to handle this:

from watchdog.events import FileSystemEventHandler
from threading import Timer
from collections import defaultdict

class DebouncedHandler(FileSystemEventHandler):
    def __init__(self, delay=0.5):
        self.delay = delay
        self._timers = {}

    def on_modified(self, event):
        if event.is_directory:
            return

        # Cancel existing timer for this file
        if event.src_path in self._timers:
            self._timers[event.src_path].cancel()

        # Set new timer
        timer = Timer(self.delay, self._process, [event.src_path])
        self._timers[event.src_path] = timer
        timer.start()

    def _process(self, path):
        print(f"Processing: {path}")
        del self._timers[path]

Pattern Matching and Filtering Events

Using PatternMatchingEventHandler

from watchdog.events import PatternMatchingEventHandler

class MyPatternHandler(PatternMatchingEventHandler):
    def __init__(self):
        super().__init__(
            patterns=["*.py", "*.json", "*.yaml"],
            ignore_patterns=["*~", "*.tmp", "__pycache__/*"],
            ignore_directories=True,
            case_sensitive=True
        )

    def on_modified(self, event):
        print(f"Matched file modified: {event.src_path}")

Using RegexMatchingEventHandler

from watchdog.events import RegexMatchingEventHandler

class MyRegexHandler(RegexMatchingEventHandler):
    def __init__(self):
        super().__init__(
            regexes=[r".*\.py$", r".*config.*\.json$"],
            ignore_regexes=[r".*test.*", r".*__pycache__.*"],
            ignore_directories=True,
            case_sensitive=False
        )

    def on_modified(self, event):
        print(f"Regex matched: {event.src_path}")

Custom Filtering

from watchdog.events import FileSystemEventHandler
import os

class FilteredHandler(FileSystemEventHandler):
    def __init__(self, extensions=None, min_size=0, exclude_hidden=True):
        self.extensions = extensions or []
        self.min_size = min_size
        self.exclude_hidden = exclude_hidden

    def _should_process(self, path):
        filename = os.path.basename(path)

        # Skip hidden files
        if self.exclude_hidden and filename.startswith('.'):
            return False

        # Check extension
        if self.extensions:
            ext = os.path.splitext(path)[1].lower()
            if ext not in self.extensions:
                return False

        # Check file size (for existing files)
        if os.path.exists(path) and os.path.getsize(path) < self.min_size:
            return False

        return True

    def on_modified(self, event):
        if not event.is_directory and self._should_process(event.src_path):
            print(f"Processing: {event.src_path}")

# Usage
handler = FilteredHandler(
    extensions=['.py', '.json'],
    min_size=100,
    exclude_hidden=True
)

Recursive vs Non-Recursive Monitoring

Recursive Monitoring

Monitors directory and all subdirectories:

observer.schedule(handler, "/path/to/watch", recursive=True)

Non-Recursive Monitoring

Monitors only the specified directory:

observer.schedule(handler, "/path/to/watch", recursive=False)

Selective Recursion

from watchdog.events import FileSystemEventHandler
import os

class SelectiveRecursionHandler(FileSystemEventHandler):
    def __init__(self, exclude_dirs=None):
        self.exclude_dirs = exclude_dirs or ['node_modules', '.git', '__pycache__']

    def on_any_event(self, event):
        # Check if event is in excluded directory
        path_parts = event.src_path.split(os.sep)
        if any(excluded in path_parts for excluded in self.exclude_dirs):
            return

        # Process event
        print(f"{event.event_type}: {event.src_path}")

# Schedule with recursive but handler filters
observer.schedule(SelectiveRecursionHandler(), "/project", recursive=True)

Multiple Non-Recursive Watches

import os

def watch_specific_dirs(base_path, subdirs, handler):
    observer = Observer()

    for subdir in subdirs:
        path = os.path.join(base_path, subdir)
        if os.path.isdir(path):
            observer.schedule(handler, path, recursive=False)

    return observer

# Watch only specific subdirectories
observer = watch_specific_dirs(
    "/project",
    ["src", "config", "templates"],
    MyHandler()
)

Async Patterns

Asyncio Integration

import asyncio
from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler
from queue import Queue

class AsyncHandler(FileSystemEventHandler):
    def __init__(self, queue):
        self.queue = queue

    def on_any_event(self, event):
        self.queue.put(event)

async def process_events(queue):
    loop = asyncio.get_event_loop()
    while True:
        # Non-blocking queue check
        event = await loop.run_in_executor(None, queue.get)
        await handle_event(event)

async def handle_event(event):
    print(f"Async processing: {event.src_path}")
    # Async file operations
    await asyncio.sleep(0.1)

async def main():
    queue = Queue()
    handler = AsyncHandler(queue)

    observer = Observer()
    observer.schedule(handler, "/path/to/watch", recursive=True)
    observer.start()

    try:
        await process_events(queue)
    finally:
        observer.stop()
        observer.join()

asyncio.run(main())

Using asyncio.Queue

import asyncio
from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler

class AsyncQueueHandler(FileSystemEventHandler):
    def __init__(self, loop, async_queue):
        self.loop = loop
        self.async_queue = async_queue

    def on_any_event(self, event):
        # Thread-safe way to put into async queue
        self.loop.call_soon_threadsafe(
            self.async_queue.put_nowait, event
        )

async def event_consumer(queue):
    while True:
        event = await queue.get()
        print(f"Processing: {event.event_type} - {event.src_path}")
        # Async processing here
        queue.task_done()

async def main():
    loop = asyncio.get_event_loop()
    queue = asyncio.Queue()

    handler = AsyncQueueHandler(loop, queue)
    observer = Observer()
    observer.schedule(handler, "/path/to/watch", recursive=True)
    observer.start()

    # Start consumer
    consumer_task = asyncio.create_task(event_consumer(queue))

    try:
        await asyncio.sleep(float('inf'))
    finally:
        observer.stop()
        consumer_task.cancel()
        observer.join()

asyncio.run(main())

Integrating with Other Python Libraries

With Logging

import logging
from watchdog.observers import Observer
from watchdog.events import LoggingEventHandler

# Configure logging
logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(message)s',
    datefmt='%Y-%m-%d %H:%M:%S'
)

# Use built-in logging handler
handler = LoggingEventHandler()
observer = Observer()
observer.schedule(handler, "/path/to/watch", recursive=True)
observer.start()

With Queue for Thread-Safe Processing

from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler
from queue import Queue
from threading import Thread

class QueueHandler(FileSystemEventHandler):
    def __init__(self, queue):
        self.queue = queue

    def on_any_event(self, event):
        self.queue.put(event)

def worker(queue):
    while True:
        event = queue.get()
        if event is None:
            break
        # Process event
        print(f"Worker processing: {event.src_path}")
        queue.task_done()

# Setup
event_queue = Queue()
worker_thread = Thread(target=worker, args=(event_queue,))
worker_thread.start()

observer = Observer()
observer.schedule(QueueHandler(event_queue), "/path", recursive=True)
observer.start()

With SQLAlchemy for Database Tracking

from watchdog.events import FileSystemEventHandler
from sqlalchemy import create_engine, Column, String, DateTime
from sqlalchemy.orm import declarative_base, sessionmaker
from datetime import datetime

Base = declarative_base()

class FileEvent(Base):
    __tablename__ = 'file_events'
    id = Column(String, primary_key=True)
    event_type = Column(String)
    path = Column(String)
    timestamp = Column(DateTime)

class DatabaseHandler(FileSystemEventHandler):
    def __init__(self, db_url):
        engine = create_engine(db_url)
        Base.metadata.create_all(engine)
        self.Session = sessionmaker(bind=engine)

    def on_any_event(self, event):
        session = self.Session()
        try:
            file_event = FileEvent(
                id=f"{event.src_path}_{datetime.now().timestamp()}",
                event_type=event.event_type,
                path=event.src_path,
                timestamp=datetime.now()
            )
            session.add(file_event)
            session.commit()
        finally:
            session.close()

With Redis for Distributed Systems

import redis
import json
from watchdog.events import FileSystemEventHandler
from datetime import datetime

class RedisHandler(FileSystemEventHandler):
    def __init__(self, redis_url='redis://localhost:6379'):
        self.redis = redis.from_url(redis_url)
        self.channel = 'file_events'

    def on_any_event(self, event):
        message = {
            'event_type': event.event_type,
            'src_path': event.src_path,
            'is_directory': event.is_directory,
            'timestamp': datetime.now().isoformat()
        }

        if hasattr(event, 'dest_path'):
            message['dest_path'] = event.dest_path

        self.redis.publish(self.channel, json.dumps(message))

Common Use Cases

Real-Time File Processing

from watchdog.events import PatternMatchingEventHandler
import subprocess
import os

class FileProcessor(PatternMatchingEventHandler):
    def __init__(self):
        super().__init__(patterns=["*.csv"])
        self.processed_dir = "/data/processed"

    def on_created(self, event):
        if not event.is_directory:
            self.process_csv(event.src_path)

    def process_csv(self, filepath):
        filename = os.path.basename(filepath)
        output = os.path.join(self.processed_dir, f"processed_{filename}")

        # Example: Run data processing script
        subprocess.run([
            "python", "process_data.py",
            "--input", filepath,
            "--output", output
        ])
        print(f"Processed: {filepath} -> {output}")

Auto-Reload Configuration

import json
import yaml
from watchdog.events import FileSystemEventHandler

class ConfigReloader(FileSystemEventHandler):
    def __init__(self, config_path, callback):
        self.config_path = config_path
        self.callback = callback
        self.config = self._load_config()

    def _load_config(self):
        with open(self.config_path) as f:
            if self.config_path.endswith('.json'):
                return json.load(f)
            elif self.config_path.endswith(('.yml', '.yaml')):
                return yaml.safe_load(f)

    def on_modified(self, event):
        if event.src_path == self.config_path:
            try:
                self.config = self._load_config()
                self.callback(self.config)
                print("Configuration reloaded successfully")
            except Exception as e:
                print(f"Failed to reload config: {e}")

# Usage
def on_config_change(new_config):
    print(f"New config: {new_config}")

handler = ConfigReloader("/etc/app/config.yaml", on_config_change)

Log File Monitoring (tail -f style)

from watchdog.events import FileSystemEventHandler
import os

class LogTailer(FileSystemEventHandler):
    def __init__(self, log_path, callback):
        self.log_path = log_path
        self.callback = callback
        self._position = os.path.getsize(log_path) if os.path.exists(log_path) else 0

    def on_modified(self, event):
        if event.src_path == self.log_path:
            self._read_new_lines()

    def _read_new_lines(self):
        with open(self.log_path, 'r') as f:
            f.seek(self._position)
            new_lines = f.readlines()
            self._position = f.tell()

            for line in new_lines:
                self.callback(line.rstrip())

# Usage
def process_log_line(line):
    if 'ERROR' in line:
        print(f"Alert: {line}")

handler = LogTailer("/var/log/app.log", process_log_line)
observer = Observer()
observer.schedule(handler, os.path.dirname("/var/log/app.log"))

Development Server Auto-Restart

from watchdog.events import PatternMatchingEventHandler
import subprocess
import sys
import os

class DevServerReloader(PatternMatchingEventHandler):
    def __init__(self, server_script):
        super().__init__(
            patterns=["*.py"],
            ignore_patterns=["*test*", "*__pycache__*"]
        )
        self.server_script = server_script
        self.process = None
        self.start_server()

    def start_server(self):
        if self.process:
            self.process.terminate()
            self.process.wait()

        self.process = subprocess.Popen([
            sys.executable, self.server_script
        ])
        print("Server started")

    def on_modified(self, event):
        if not event.is_directory:
            print(f"Detected change in {event.src_path}")
            self.start_server()

    def cleanup(self):
        if self.process:
            self.process.terminate()

Backup System

from watchdog.events import FileSystemEventHandler
import shutil
import os
from datetime import datetime

class BackupHandler(FileSystemEventHandler):
    def __init__(self, backup_dir):
        self.backup_dir = backup_dir
        os.makedirs(backup_dir, exist_ok=True)

    def on_modified(self, event):
        if not event.is_directory:
            self._create_backup(event.src_path)

    def on_created(self, event):
        if not event.is_directory:
            self._create_backup(event.src_path)

    def _create_backup(self, filepath):
        filename = os.path.basename(filepath)
        timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
        backup_name = f"{timestamp}_{filename}"
        backup_path = os.path.join(self.backup_dir, backup_name)

        try:
            shutil.copy2(filepath, backup_path)
            print(f"Backup created: {backup_path}")
        except Exception as e:
            print(f"Backup failed: {e}")

Performance Considerations

Optimising for High-Volume Directories

from watchdog.events import FileSystemEventHandler
from threading import Lock
from collections import deque
import time

class BatchingHandler(FileSystemEventHandler):
    def __init__(self, batch_size=100, flush_interval=1.0):
        self.batch_size = batch_size
        self.flush_interval = flush_interval
        self.events = deque()
        self.lock = Lock()
        self.last_flush = time.time()

    def on_any_event(self, event):
        with self.lock:
            self.events.append(event)

            # Flush if batch size reached or interval passed
            if (len(self.events) >= self.batch_size or
                time.time() - self.last_flush >= self.flush_interval):
                self._flush()

    def _flush(self):
        events_to_process = list(self.events)
        self.events.clear()
        self.last_flush = time.time()

        # Process batch
        self._process_batch(events_to_process)

    def _process_batch(self, events):
        print(f"Processing batch of {len(events)} events")
        # Batch processing logic here

Memory Management

from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler
import gc

class MemoryEfficientHandler(FileSystemEventHandler):
    def __init__(self):
        self.event_count = 0
        self.gc_threshold = 10000

    def on_any_event(self, event):
        self.event_count += 1

        # Process event without storing references
        self._process(event.event_type, event.src_path)

        # Periodic garbage collection
        if self.event_count % self.gc_threshold == 0:
            gc.collect()

    def _process(self, event_type, path):
        # Process with minimal memory footprint
        pass

Platform-Specific Optimisations

import sys
from watchdog.observers import Observer

def create_optimised_observer():
    """Create observer with platform-specific settings."""
    if sys.platform == 'linux':
        # Linux: Use polling for network filesystems
        from watchdog.observers.polling import PollingObserver
        return PollingObserver(timeout=1)
    elif sys.platform == 'darwin':
        # macOS: FSEvents is efficient
        return Observer()
    elif sys.platform == 'win32':
        # Windows: ReadDirectoryChangesW
        return Observer()
    else:
        # Fallback to polling
        from watchdog.observers.polling import PollingObserver
        return PollingObserver()

Reducing Event Noise

from watchdog.events import FileSystemEventHandler
import os

class NoiseReducingHandler(FileSystemEventHandler):
    # Files that frequently cause noise
    NOISE_PATTERNS = {
        '.swp', '.swx', '~', '.tmp', '.temp',
        '.DS_Store', 'Thumbs.db', '.git'
    }

    def _is_noise(self, path):
        basename = os.path.basename(path)
        return any(
            basename.endswith(pattern) or pattern in path
            for pattern in self.NOISE_PATTERNS
        )

    def on_any_event(self, event):
        if self._is_noise(event.src_path):
            return

        # Process legitimate event
        print(f"{event.event_type}: {event.src_path}")

Quick Reference

Essential Imports

from watchdog.observers import Observer
from watchdog.events import (
    FileSystemEventHandler,
    PatternMatchingEventHandler,
    RegexMatchingEventHandler,
    LoggingEventHandler
)

Minimal Working Example

from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler
import time

class Handler(FileSystemEventHandler):
    def on_modified(self, event):
        print(f"Modified: {event.src_path}")

observer = Observer()
observer.schedule(Handler(), ".", recursive=True)
observer.start()

try:
    while True:
        time.sleep(1)
except KeyboardInterrupt:
    observer.stop()
observer.join()

Event Handler Methods

Method Triggered When
on_any_event(event) Any event occurs
on_created(event) File/directory created
on_modified(event) File/directory modified
on_deleted(event) File/directory deleted
on_moved(event) File/directory moved
on_closed(event) File closed (Linux only)

Observer Methods

Method Description
schedule(handler, path, recursive) Register handler for path
unschedule(watch) Stop watching a path
unschedule_all() Stop all watches
start() Start observer thread
stop() Stop observer thread
join() Wait for observer to finish
is_alive() Check if observer running

Pattern Handler Parameters

PatternMatchingEventHandler(
    patterns=["*.py"],           # Patterns to match
    ignore_patterns=["*~"],      # Patterns to ignore
    ignore_directories=False,    # Ignore directory events
    case_sensitive=True          # Case-sensitive matching
)

Common Issues and Solutions

Issue: Multiple Events for Single File Save

Problem: Editors trigger multiple modify events when saving.

Solution: Implement debouncing:

from threading import Timer

class DebouncedHandler(FileSystemEventHandler):
    def __init__(self, delay=0.5):
        self.delay = delay
        self._timers = {}

    def on_modified(self, event):
        if event.src_path in self._timers:
            self._timers[event.src_path].cancel()

        timer = Timer(self.delay, self._handle, [event])
        self._timers[event.src_path] = timer
        timer.start()

    def _handle(self, event):
        del self._timers[event.src_path]
        # Process event here

Issue: Events Not Detected on Network Drives

Problem: inotify doesn't work on NFS/CIFS mounts.

Solution: Use polling observer:

from watchdog.observers.polling import PollingObserver

observer = PollingObserver(timeout=2)  # Poll every 2 seconds
observer.schedule(handler, "/mnt/network/share", recursive=True)

Issue: Permission Denied Errors

Problem: Cannot watch directories without read permission.

Solution: Handle exceptions gracefully:

import os

def safe_schedule(observer, handler, path, recursive=True):
    if not os.path.exists(path):
        print(f"Path does not exist: {path}")
        return None

    if not os.access(path, os.R_OK):
        print(f"No read permission: {path}")
        return None

    try:
        return observer.schedule(handler, path, recursive=recursive)
    except Exception as e:
        print(f"Failed to watch {path}: {e}")
        return None

Issue: High CPU Usage

Problem: Polling observer or high event volume causes CPU spikes.

Solutions:

  1. Increase polling interval:
PollingObserver(timeout=5)  # 5 seconds instead of default 1
  1. Filter events early:
def on_any_event(self, event):
    if event.is_directory:
        return  # Skip directory events
    if not event.src_path.endswith('.py'):
        return  # Skip non-Python files
  1. Batch process events (see Performance Considerations section)

Issue: Observer Stops Without Error

Problem: Observer thread dies silently.

Solution: Add error handling and monitoring:

import threading
import time

def monitor_observer(observer):
    while True:
        if not observer.is_alive():
            print("Observer died, restarting...")
            observer.start()
        time.sleep(5)

# Start monitor thread
monitor = threading.Thread(target=monitor_observer, args=(observer,))
monitor.daemon = True
monitor.start()

Issue: Too Many Open Files

Problem: OSError: [Errno 24] Too many open files when watching many directories.

Solution: Increase system limits or reduce watch scope:

# Check current limit
ulimit -n

# Increase limit (temporary)
ulimit -n 65535

# Permanent: edit /etc/security/limits.conf
# * soft nofile 65535
# * hard nofile 65535

Or use non-recursive watching with specific subdirectories.

Issue: Missed Events During Heavy Load

Problem: Events dropped when filesystem is busy.

Solution: Use queuing to prevent blocking:

from queue import Queue
from threading import Thread

class QueuedHandler(FileSystemEventHandler):
    def __init__(self):
        self.queue = Queue(maxsize=10000)
        self._start_worker()

    def _start_worker(self):
        def worker():
            while True:
                event = self.queue.get()
                self._process(event)
                self.queue.task_done()

        thread = Thread(target=worker, daemon=True)
        thread.start()

    def on_any_event(self, event):
        try:
            self.queue.put_nowait(event)
        except:
            print("Queue full, event dropped")

    def _process(self, event):
        # Actual processing here
        pass

Issue: File Not Ready When Event Received

Problem: File is still being written when created/modified event fires.

Solution: Wait for file to be ready:

import os
import time

def wait_for_file(path, timeout=30):
    """Wait for file to stop being written to."""
    last_size = -1
    stable_count = 0

    start = time.time()
    while time.time() - start < timeout:
        try:
            current_size = os.path.getsize(path)
            if current_size == last_size:
                stable_count += 1
                if stable_count >= 3:
                    return True
            else:
                stable_count = 0
                last_size = current_size
        except OSError:
            pass
        time.sleep(0.5)

    return False

def on_created(self, event):
    if wait_for_file(event.src_path):
        self.process_file(event.src_path)