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

如何将Watchdog与FastAPI集成实现文件夹监控并保存路径到数据库

FastAPI集成Watchdog实现文件监控并存储路径到数据库

问题说明

需要实现文件夹新增文件时自动将文件路径存入数据库,选用watchdog实现文件系统监控,但不清楚如何与FastAPI正确集成。已尝试修改代码实现FastAPI启动时自动开启监控,但存在依赖注入的错误。

原代码问题点

原代码中MyHandler的__init__方法直接使用Depends(get_session)是错误的——Depends是FastAPI专为路由函数、依赖项设计的注入工具,无法在普通类的初始化方法中生效,会导致无法正确获取数据库会话。

修正后的完整实现

1. 核心代码(FastAPI + Watchdog)

from fastapi import FastAPI
from sqlalchemy.orm import Session
from session import get_session  # 你的数据库会话生成器
from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler
import threading

app = FastAPI()

# 示例:文件路径存储模型(需根据你的实际数据库结构调整)
# 假设session.py中已定义Base和数据库连接
class FileRecord(Base):
    __tablename__ = "file_records"
    path = Column(String, primary_key=True, index=True)

class MyHandler(FileSystemEventHandler):
    def __init__(self, target_folder, db_session_generator):
        super().__init__()
        self.target_folder = target_folder
        self.get_db = db_session_generator

    def on_created(self, event):
        # 忽略文件夹,只处理文件
        if event.is_directory:
            return
        
        file_path = event.src_path
        print(f"检测到新文件: {file_path}")
        
        # 获取数据库会话并存储路径
        db: Session = next(self.get_db())
        try:
            # 避免重复存储同一文件路径
            if not db.query(FileRecord).filter(FileRecord.path == file_path).first():
                new_record = FileRecord(path=file_path)
                db.add(new_record)
                db.commit()
                print(f"文件路径已存入数据库: {file_path}")
        except Exception as e:
            db.rollback()
            print(f"存储失败: {str(e)}")
        finally:
            db.close()

    def on_modified(self, event):
        if not event.is_directory:
            print(f"文件已修改: {event.src_path}")

class WatchdogThread:
    def __init__(self, target_folder, db_session_generator):
        self.observer = Observer()
        self.target_folder = target_folder
        self.get_db = db_session_generator

    def start(self):
        event_handler = MyHandler(self.target_folder, self.get_db)
        self.observer.schedule(event_handler, self.target_folder, recursive=True)
        print(f"启动文件夹监控: {self.target_folder}")
        # 守护线程启动,避免阻塞FastAPI主线程
        threading.Thread(target=self.observer.start, daemon=True).start()

    def stop(self):
        self.observer.stop()
        print(f"停止文件夹监控: {self.target_folder}")
        self.observer.join()

# 监控目标路径(与Docker挂载路径一致)
target_folder = "/monitored"
watchdog_thread = WatchdogThread(target_folder, get_session)

# FastAPI启动/停止钩子
@app.on_event("startup")
async def startup():
    watchdog_thread.start()

@app.on_event("shutdown")
async def shutdown():
    watchdog_thread.stop()

def main():
    import uvicorn
    uvicorn.run("app.main:app", host="0.0.0.0", port=8000, access_log=True, reload=True)

if __name__ == "__main__":
    main()

2. Docker Compose配置

version: '3.8'
services:
  fastapi-watchdog:
    build: .
    volumes:
      # 将宿主机目标文件夹挂载到容器内的监控路径
      - ${TARGET_FOLDER_PARENT_PATH}:/monitored:rw
    ports:
      - "8000:8000"
    environment:
      - TARGET_FOLDER_PARENT_PATH=${TARGET_FOLDER_PARENT_PATH}

关键注意事项

  • 依赖注入修正:通过传入数据库会话生成器get_session,在事件处理函数中主动获取会话,替代原错误的Depends用法。
  • 线程隔离:用守护线程启动Watchdog的Observer,避免阻塞FastAPI的启动和运行流程。
  • 数据库操作安全:添加异常捕获、事务回滚和会话关闭逻辑,确保数据库操作的稳定性。
  • 重复存储规避:新增文件时先查询数据库,避免重复插入同一文件路径。

内容的提问来源于stack exchange,提问作者hfc

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 11:54:58