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

求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>

四、完整运行流程

  1. 启动Redis服务器(如果还没运行的话)
  2. 启动Celery Worker:
    celery -A your_task_module worker --loglevel=info
    
  3. 启动Celery的Watchdog监听任务:
    from your_task_module import start_watcher
    start_watcher.delay()
    
  4. 启动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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:06:55