求Flask集成Watchdog Observer示例及Celery部署与前端更新优化方案
集成Flask、Watchdog和Celery的完整方案
看起来你已经走了大半路程——用Celery跑Watchdog监听、Handler更新数据库都搞定了,剩下的就是让前端更优雅地接收变化通知,以及优化Flask和Celery之间的事件同步。下面是一步步的优化方案:
一、优化Celery中的Watchdog监听逻辑
你的现有代码里用了event_q和循环检查,其实可以简化:让MyHandler在触发事件时直接完成两个动作——更新数据库 + 把事件信息发送到消息中间件(比如Redis,毕竟Celery大概率已经用Redis/RabbitMQ做 broker 了),这样Celery任务不需要一直轮询队列,专注维持Watchdog运行即可。
调整后的Celery任务和Handler示例:
import json import time import redis from celery import Celery from watchdog.observers import Observer from watchdog.events import FileSystemEventHandler celery = Celery('watcher_tasks', broker='redis://localhost:6379/0') redis_client = redis.Redis(host='localhost', port=6379, db=0) class MyHandler(FileSystemEventHandler): def on_any_event(self, event): if not event.is_directory: # 按需过滤目录事件 # 1. 已实现:更新数据库 self.update_database(event) # 2. 新增:把事件推送到Redis的消息频道 event_data = { 'src_path': event.src_path, 'event_type': event.event_type, 'timestamp': event.timestamp } redis_client.publish('file_changes', json.dumps(event_data)) def update_database(self, event): # 这里放你已有的数据库更新逻辑 pass @celery.task(bind=True) def start_watcher(self): observer = Observer() handler = MyHandler() observer.schedule(handler, '.', recursive=True) # 显式开启子目录监听 observer.start() try: # 保持任务运行,无需轮询event_q while True: time.sleep(1) except KeyboardInterrupt: observer.stop() observer.join()
二、Flask后端实时转发事件到前端
用Flask-SocketIO替代前端轮询,让Flask主动推送变化给前端。Flask端需要监听Redis的消息频道,一旦收到事件就通过SocketIO推送给所有连接的客户端。
首先安装依赖:
pip install flask-socketio redis eventlet
Flask核心代码示例:
from flask import Flask, render_template from flask_socketio import SocketIO, emit import redis import json app = Flask(__name__) app.config['SECRET_KEY'] = 'your-secret-key-here' socketio = SocketIO(app, cors_allowed_origins="*") # 按需调整跨域配置 # 后台线程监听Redis消息频道 def listen_redis_changes(): redis_client = redis.Redis(host='localhost', port=6379, db=0) pubsub = redis_client.pubsub() pubsub.subscribe('file_changes') for message in pubsub.listen(): if message['type'] == 'message': event_data = json.loads(message['data']) # 通过SocketIO推送给所有前端客户端 socketio.emit('file_change', event_data, namespace='/') # 启动后台监听线程 import threading thread = threading.Thread(target=listen_redis_changes) thread.daemon = True thread.start() @app.route('/') def index(): return render_template('index.html') # 你的前端页面模板 if __name__ == '__main__': socketio.run(app, debug=True)
三、前端实时接收更新
前端用SocketIO客户端监听file_change事件,不需要每秒轮询接口,直接接收后端推送的变化并更新页面。
前端HTML/JS示例(放在templates/index.html):
<!DOCTYPE html> <html> <head> <title>File System Watcher</title> <script src="https://cdnjs.cloudflare.com/ajax/libs/socket.io/4.0.1/socket.io.js"></script> </head> <body> <h1>Latest File Changes</h1> <ul id="changes-list"></ul> <script> const socket = io(); const changesList = document.getElementById('changes-list'); socket.on('file_change', function(data) { // 创建新列表项展示变化 const li = document.createElement('li'); li.textContent = `${data.event_type}: ${data.src_path} (${new Date(data.timestamp * 1000).toLocaleString()})`; changesList.prepend(li); // 最新变化放在最顶部 }); </script> </body> </html>
四、完整运行流程
- 启动Redis服务器(如果还没运行的话)
- 启动Celery Worker:
celery -A your_task_module worker --loglevel=info - 启动Celery的Watchdog监听任务:
from your_task_module import start_watcher start_watcher.delay() - 启动Flask应用:
python app.py
关键说明
- 为什么用Redis做中间件?:Celery任务和Flask是独立进程,Redis的发布/订阅功能能让它们高效解耦通信,避免Celery直接依赖Flask上下文,也保证了事件的可靠性。
- 事件过滤优化:可以在
MyHandler中单独实现on_created/on_modified/on_deleted等方法,按需过滤不需要的事件(比如临时文件、隐藏文件)。 - 替代方案:如果不想用Redis,也可以在Celery任务中直接调用Flask的SocketIO客户端,但需要确保Flask应用地址可访问,这种方式解耦性不如Redis好。
内容的提问来源于stack exchange,提问作者Sungsoo Ha
相关产品推荐
相关产品推荐

