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

如何用队列存储批量新增文件,解决FileSystemWatcher批量处理问题

批量文件监控处理的队列实现方案

问题背景

开发了一个监控文件夹的服务,当有新文件加入时会逐行解析内容并将每行数据编码为二维码保存到新文件夹。单个文件加入时服务可正常运行,但批量放入多个文件时仅能处理其中一个。需要实现一个队列来存储所有检测到的新增文件,完成当前文件编码后自动取出下一个文件处理,处理完成后删除原文件。计划使用System.Collections.Concurrent队列实现。

现有代码

启动服务/文件系统监控

public void Start()
{
    newfolders = ConfigurationManager.AppSettings["parsefiles"];
    pathtomonitor = ConfigurationManager.AppSettings["pathtomonitor"];

    pathforpaths = newfolders;

    System.IO.Directory.CreateDirectory(pathforpaths);
    Thread.Sleep(2000);

    pathforprocess = newfolders + @"\processdata";

    System.IO.Directory.CreateDirectory(pathtomonitor);
    System.IO.Directory.CreateDirectory(pathforprocess);

    FileSystemWatcher watcher = new FileSystemWatcher(pathtomonitor);

    watcher.NotifyFilter = NotifyFilters.Attributes | NotifyFilters.FileName;

    watcher.EnableRaisingEvents = true;
    watcher.IncludeSubdirectories = false;

    //add event handlers
    watcher.Created += watcher_Created;
}

原文件处理逻辑

private static void watcher_Created(object sender, FileSystemEventArgs e)
{
    if (e.ChangeType == WatcherChangeTypes.Created)
    {

        if (DateTime.Now.Subtract(_lastTimeFileWatcherEventRaised)
            .TotalMilliseconds < 1000)
        {
            Console.WriteLine("Short time between previous event and change event");
            return;
        }

        Console.WriteLine("file: {0} changed at time:{1}", e.Name,
            DateTime.Now.ToLocalTime());

        Thread.Sleep(1000);
        if (GetAccessLoop(pathtomonitor + "\\" + e.Name))
        {
            parsefile(pathtomonitor + "\\" + e.Name);
            Console.WriteLine("worked");
        }
        else
        {
            //writetologs("Was not able to get access to" + e.Name);
            Console.WriteLine("didnt work");
        }

        _lastTimeFileWatcherEventRaised = DateTime.Now;
    }
}

解决方案

核心思路

使用**线程安全的ConcurrentQueue<string>**存储待处理的文件路径,让FileSystemWatcher的事件处理器仅负责将文件路径加入队列,启动一个独立的后台线程专门从队列中取出文件并处理,实现解耦和批量文件的顺序处理。不需要使用MemoryStream,因为我们只需要存储文件路径,无需提前加载文件内容。

具体实现步骤

  1. 添加全局队列和处理线程字段
    在类中定义并发队列和处理线程的成员变量,同时添加取消令牌用于优雅停止线程:
private readonly ConcurrentQueue<string> _fileQueue = new ConcurrentQueue<string>();
private Thread _processingThread;
private CancellationTokenSource _cts;
private static DateTime _lastTimeFileWatcherEventRaised;
// 保留原有的路径字段:newfolders, pathtomonitor, pathforpaths, pathforprocess
  1. 修改Start方法,启动处理线程
    在Start方法末尾添加启动后台处理线程的逻辑:
public void Start()
{
    // 原有的初始化代码...

    FileSystemWatcher watcher = new FileSystemWatcher(pathtomonitor);
    // 原有的watcher配置代码...
    watcher.Created += watcher_Created;

    // 启动处理队列的后台线程
    _cts = new CancellationTokenSource();
    _processingThread = new Thread(() => ProcessQueue(_cts.Token))
    {
        IsBackground = true // 设置为后台线程,避免阻止程序退出
    };
    _processingThread.Start();
}
  1. 修改Watcher事件处理器,仅负责入队
    将原有的处理逻辑简化为将文件路径加入队列,避免事件阻塞:
private static void watcher_Created(object sender, FileSystemEventArgs e)
{
    if (e.ChangeType != WatcherChangeTypes.Created) return;

    // 可选:保留短时间间隔过滤,避免重复触发
    if (DateTime.Now.Subtract(_lastTimeFileWatcherEventRaised).TotalMilliseconds < 1000)
    {
        Console.WriteLine("Short time between previous event and change event");
        return;
    }

    string filePath = Path.Combine(pathtomonitor, e.Name);
    Console.WriteLine($"File added to queue: {filePath} at {DateTime.Now.ToLocalTime()}");
    _fileQueue.Enqueue(filePath);

    _lastTimeFileWatcherEventRaised = DateTime.Now;
}
  1. 实现队列处理循环
    编写后台线程的处理逻辑,循环从队列取文件并处理,处理完成后删除原文件:
private void ProcessQueue(CancellationToken token)
{
    while (!token.IsCancellationRequested)
    {
        if (_fileQueue.TryDequeue(out string filePath))
        {
            try
            {
                Console.WriteLine($"Starting processing: {filePath}");
                // 等待文件可访问(保留原有的GetAccessLoop逻辑)
                if (GetAccessLoop(filePath))
                {
                    parsefile(filePath);
                    Console.WriteLine($"Processed successfully: {filePath}");
                    // 处理完成后删除文件
                    File.Delete(filePath);
                    Console.WriteLine($"Deleted processed file: {filePath}");
                }
                else
                {
                    Console.WriteLine($"Failed to get access to file: {filePath}, re-enqueuing...");
                    // 访问失败时重新加入队列,稍后重试
                    _fileQueue.Enqueue(filePath);
                }
            }
            catch (Exception ex)
            {
                Console.WriteLine($"Error processing file {filePath}: {ex.Message}");
                // 可选:将失败文件重新入队或记录日志
                _fileQueue.Enqueue(filePath);
            }
        }
        else
        {
            // 队列为空时短暂休眠,避免CPU空转
            Thread.Sleep(500);
        }
    }
}
  1. 添加停止服务的方法(可选)
    如果需要优雅停止服务,添加取消令牌的触发逻辑:
public void Stop()
{
    _cts?.Cancel();
    _processingThread?.Join();
}

关键注意事项

  • 线程安全:ConcurrentQueue是线程安全的,无需额外加锁即可在多线程环境下(Watcher事件线程和处理线程)操作。
  • 文件访问问题:保留GetAccessLoop确保文件写入完成后再处理,避免因文件被占用导致失败。
  • 异常处理:在处理循环中捕获异常,避免单个文件处理失败导致整个处理线程崩溃,同时可将失败文件重新入队重试。
  • 避免CPU空转:队列为空时短暂休眠,减少资源消耗。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 16:45:17