FastAPI场景下如何实现watchdog联动WebSocket推送目录文件变更?
解决方案:Watchdog 与 FastAPI WebSocket 的实时通信方案
针对你提出的「Watchdog监控脚本如何向FastAPI的WebSocket客户端推送文件变更」问题,以下是三种通用最佳实践方案,适配不同场景需求:
方案1:将Watchdog集成到FastAPI异步服务中(推荐轻量场景)
原理:利用ASGI的异步事件循环,把Watchdog监控任务作为后台任务嵌入FastAPI服务,直接共享WebSocket连接池,无需跨进程通信,部署简单。
实现代码
import asyncio from fastapi import FastAPI, WebSocket, WebSocketDisconnect from watchdog.observers.asyncio import AsyncObserver from watchdog.events import FileSystemEventHandler app = FastAPI() # 存储活跃WebSocket连接,用锁保证异步安全 active_connections = set() conn_lock = asyncio.Lock() class PhotoChangeHandler(FileSystemEventHandler): async def on_created(self, event): if not event.is_directory and event.src_path.lower().endswith(('.jpg', '.png')): await push_to_clients({"type": "add", "path": event.src_path}) async def on_deleted(self, event): if not event.is_directory and event.src_path.lower().endswith(('.jpg', '.png')): await push_to_clients({"type": "remove", "path": event.src_path}) async def push_to_clients(message): async with conn_lock: for connection in active_connections: await connection.send_json(message) @app.websocket("/ws/photos") async def websocket_endpoint(websocket: WebSocket): await websocket.accept() async with conn_lock: active_connections.add(websocket) try: while True: # 维持连接,无需接收客户端消息可省略此句 await websocket.receive_text() except WebSocketDisconnect: async with conn_lock: active_connections.remove(websocket) @app.on_event("startup") async def startup_event(): # 初始化异步Watchdog观察者 event_handler = PhotoChangeHandler() observer = AsyncObserver() observer.schedule(event_handler, path="/path/to/photos", recursive=False) observer.start() app.state.observer = observer @app.on_event("shutdown") async def shutdown_event(): app.state.observer.stop() await app.state.observer.join()
适用场景:轻量部署,无需拆分Watchdog进程,Docker中直接打包成单个镜像即可运行。
方案2:用Redis Pub/Sub实现解耦通信(适合扩展场景)
原理:引入Redis作为消息中间件,Watchdog捕获事件后发布到Redis频道,FastAPI的WebSocket服务订阅该频道,收到消息后推送给客户端。完全解耦两个组件,支持后续多实例扩展。
实现代码
FastAPI WebSocket服务
import asyncio from fastapi import FastAPI, WebSocket, WebSocketDisconnect import aioredis app = FastAPI() active_connections = set() conn_lock = asyncio.Lock() REDIS_CHANNEL = "photo_changes" async def redis_listener(): redis = await aioredis.from_url("redis://localhost") pubsub = redis.pubsub() await pubsub.subscribe(REDIS_CHANNEL) async for message in pubsub.listen(): if message["type"] == "message": await push_to_clients(message["data"].decode()) async def push_to_clients(message): async with conn_lock: for connection in active_connections: await connection.send_text(message) @app.websocket("/ws/photos") async def websocket_endpoint(websocket: WebSocket): await websocket.accept() async with conn_lock: active_connections.add(websocket) try: while True: await websocket.receive_text() except WebSocketDisconnect: async with conn_lock: active_connections.remove(websocket) @app.on_event("startup") async def startup_event(): # 启动Redis监听后台任务 asyncio.create_task(redis_listener())
Watchdog监控脚本
import redis from watchdog.observers import Observer from watchdog.events import FileSystemEventHandler REDIS_CHANNEL = "photo_changes" r = redis.Redis(host='localhost', port=6379, db=0) class PhotoChangeHandler(FileSystemEventHandler): def on_created(self, event): if not event.is_directory and event.src_path.lower().endswith(('.jpg', '.png')): r.publish(REDIS_CHANNEL, f'{{"type": "add", "path": "{event.src_path}"}}') def on_deleted(self, event): if not event.is_directory and event.src_path.lower().endswith(('.jpg', '.png')): r.publish(REDIS_CHANNEL, f'{{"type": "remove", "path": "{event.src_path}"}}') if __name__ == "__main__": event_handler = PhotoChangeHandler() observer = Observer() observer.schedule(event_handler, path="/path/to/photos", recursive=False) observer.start() try: while True: pass except KeyboardInterrupt: observer.stop() observer.join()
适用场景:需要拆分服务、后续可能扩展多节点的场景,Docker部署时用docker-compose管理Redis、FastAPI、Watchdog三个服务。
方案3:内部API调用(最简但耦合度高)
原理:FastAPI暴露一个内部POST接口,Watchdog捕获事件后调用该接口,FastAPI收到请求后推送给WebSocket客户端。实现最简单,但两者强耦合,一方故障会影响另一方。
实现代码
FastAPI服务
import asyncio from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Request app = FastAPI() active_connections = set() conn_lock = asyncio.Lock() async def push_to_clients(message): async with conn_lock: for connection in active_connections: await connection.send_json(message) @app.websocket("/ws/photos") async def websocket_endpoint(websocket: WebSocket): await websocket.accept() async with conn_lock: active_connections.add(websocket) try: while True: await websocket.receive_text() except WebSocketDisconnect: async with conn_lock: active_connections.remove(websocket) @app.post("/internal/photo-change") async def receive_photo_change(request: Request): change_data = await request.json() await push_to_clients(change_data) return {"status": "ok"}
Watchdog脚本
import httpx from watchdog.observers import Observer from watchdog.events import FileSystemEventHandler API_URL = "http://localhost:8000/internal/photo-change" class PhotoChangeHandler(FileSystemEventHandler): def on_created(self, event): if not event.is_directory and event.src_path.lower().endswith(('.jpg', '.png')): httpx.post(API_URL, json={"type": "add", "path": event.src_path}) def on_deleted(self, event): if not event.is_directory and event.src_path.lower().endswith(('.jpg', '.png')): httpx.post(API_URL, json={"type": "remove", "path": event.src_path}) if __name__ == "__main__": event_handler = PhotoChangeHandler() observer = Observer() observer.schedule(event_handler, path="/path/to/photos", recursive=False) observer.start() try: while True: pass except KeyboardInterrupt: observer.stop() observer.join()
适用场景:快速原型开发,无需考虑扩展和高可用性的场景。
关键注意事项
- 异步安全:操作WebSocket连接集合时必须用
asyncio.Lock加锁,避免并发修改导致异常 - Docker部署:
- 方案1:直接打包成单个镜像,启动命令用
uvicorn main:app --host 0.0.0.0 --port 8000 - 方案2/3:用
docker-compose管理多服务,确保容器间网络互通
- 方案1:直接打包成单个镜像,启动命令用
- 事件过滤:Watchdog会捕获所有文件变更,需过滤目录和非照片文件,避免无效推送
内容的提问来源于stack exchange,提问作者Mitchell Currie
相关产品推荐
相关产品推荐

