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.
flowchart LR
subgraph "File System"
A[Files/Directories]
end
subgraph "Watchdog"
B[Observer] --> C[Event Handler]
C --> D[on_created]
C --> E[on_modified]
C --> F[on_deleted]
C --> G[on_moved]
end
A -->|"Events"| B
D --> H[Your Code]
E --> H
F --> H
G --> H
Installation
pip install watchdog
Setting Up Observers and Event Handlers
Observers monitor directories whilst event handlers process detected changes.
sequenceDiagram
participant O as Observer
participant S as Scheduler
participant H as EventHandler
participant FS as File System
O->>S: schedule(handler, path)
O->>O: start()
FS-->>S: File event detected
S->>H: dispatch(event)
H->>H: 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:
- Increase polling interval:
PollingObserver(timeout=5) # 5 seconds instead of default 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
- 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)