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

FastAPI WebSocket结合Watchdog监听文件变化无消息发送求助

问题分析与修复方案

核心问题:Watchdog事件处理器方法不能是异步的

Watchdog的FileSystemEventHandler中的回调方法(如on_modified)必须是同步函数,因为观察者线程是同步执行这些方法的。你把on_modified定义为async def,导致Watchdog无法正确触发这个方法,自然不会发送消息。

修复步骤

1. 修正事件处理器的on_modified方法,通过事件循环调度异步发送

将on_modified改为同步函数,使用asyncio把WebSocket的异步发送任务提交到FastAPI的事件循环中执行:

import asyncio
from watchdog.events import FileSystemEventHandler
from fastapi import WebSocket
import os

class LogFileEventHandler(FileSystemEventHandler):
    def __init__(self, websocket: WebSocket, videoId: str, temp_dir: str):
        self.websocket = websocket
        self.videoId = videoId
        self.temp_dir = temp_dir
        self.loop = asyncio.get_event_loop()
        # 记录文件最后修改时间,避免重复触发
        self.last_modified = 0

    def on_modified(self, event):
        # 提取文件名,避免路径前缀影响判断
        filename = os.path.basename(event.src_path)
        target_filename = f"log_{self.videoId}.txt"
        if filename == target_filename and not event.is_directory:
            log_file = os.path.join(self.temp_dir, target_filename)
            try:
                current_modified = os.path.getmtime(log_file)
                # 仅当文件真正更新时发送内容
                if current_modified > self.last_modified:
                    self.last_modified = current_modified
                    with open(log_file, "r") as file:
                        content = file.read()
                    # 提交异步发送任务到事件循环
                    self.loop.create_task(self.websocket.send_text(content))
            except Exception as e:
                print(f"处理日志文件出错: {e}")

2. 完善WebSocket端点的资源清理逻辑

确保在任何退出场景下都停止观察者,避免线程泄漏:

@router.websocket("/logs")
async def websocket_endpoint(websocket: WebSocket, videoId: str = Query(...)):
    await manager.connect(websocket)
    print(f"Client connected: {websocket.client.host}")
 
    event_handler = LogFileEventHandler(websocket, videoId, temp_dir)
    observer = Observer()
    observer.schedule(event_handler, temp_dir, recursive=False)
    observer.start()

    try:
        while True:
            try:
                # 若不需要接收前端消息,可替换为 await asyncio.sleep(3600) 维持连接
                data = await websocket.receive_json()
            except RuntimeError:
                break  
    except WebSocketDisconnect:
        print(f"Client disconnected: {websocket.client.host}") 
    finally:
        # 无论何种退出方式,都清理资源
        manager.disconnect(websocket)
        observer.stop()
        observer.join()

3. 额外检查点

  • 确认temp_dir是绝对路径,避免相对路径导致监听目录错误
  • 验证目标日志文件确实生成在temp_dir下,且文件名与log_{videoId}.txt完全匹配

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 10:57:15