如何在C#中结合FileSystemWatcher实现FIFO处理?
在C#中结合FileSystemWatcher实现FIFO文件处理
要实现按先进先出顺序处理监控目录的文件,核心是用线程安全的队列维护待处理文件,配合后台线程逐个消费队列任务。FileSystemWatcher的事件是异步触发的,直接在事件回调中处理文件会导致并发问题和顺序混乱,队列能严格保证FIFO顺序。
实现核心思路
- 用
ConcurrentQueue<string>存储待处理文件路径(线程安全,无需额外锁) - 配置
FileSystemWatcher监控目标目录的文件创建/重命名事件 - 启动后台任务持续从队列取文件处理,处理前确保文件已释放(避免粘贴大文件时文件仍被占用)
完整示例代码
using System; using System.Collections.Concurrent; using System.IO; using System.Threading; using System.Threading.Tasks; public class FifoFileProcessor { private readonly ConcurrentQueue<string> _fileQueue = new ConcurrentQueue<string>(); private readonly FileSystemWatcher _watcher; private readonly string _monitorDirectory; private CancellationTokenSource _cts; private Task _processingTask; public FifoFileProcessor(string monitorDirectory) { _monitorDirectory = monitorDirectory ?? throw new ArgumentNullException(nameof(monitorDirectory)); // 初始化FileSystemWatcher _watcher = new FileSystemWatcher { Path = _monitorDirectory, NotifyFilter = NotifyFilters.FileName | NotifyFilters.CreationTime, Filter = "*.*", // 可修改为特定后缀,比如"*.txt" EnableRaisingEvents = false // 初始化完成后再启用事件 }; // 绑定文件创建和重命名事件(覆盖剪切到目录的场景) _watcher.Created += OnFileCreated; _watcher.Renamed += OnFileRenamed; } // 启动监控与处理流程 public void StartProcessing() { if (_processingTask != null && !_processingTask.IsCompleted) throw new InvalidOperationException("Processing is already running."); _cts = new CancellationTokenSource(); // 启动长后台任务,避免占用线程池资源 _processingTask = Task.Run(async () => await ProcessFilesAsync(_cts.Token), _cts.Token); _watcher.EnableRaisingEvents = true; Console.WriteLine($"Started monitoring: {_monitorDirectory}"); } // 停止监控与处理 public async Task StopProcessingAsync() { _watcher.EnableRaisingEvents = false; _cts.Cancel(); try { await _processingTask; } catch (OperationCanceledException) { // 任务被取消属于正常流程 } finally { _watcher.Dispose(); _cts.Dispose(); Console.WriteLine("Processing stopped."); } } private void OnFileCreated(object sender, FileSystemEventArgs e) { _fileQueue.Enqueue(e.FullPath); Console.WriteLine($"Added to queue: {e.Name}"); } private void OnFileRenamed(object sender, RenamedEventArgs e) { _fileQueue.Enqueue(e.FullPath); Console.WriteLine($"Added renamed file to queue: {e.Name}"); } private async Task ProcessFilesAsync(CancellationToken cancellationToken) { while (!cancellationToken.IsCancellationRequested) { if (_fileQueue.TryDequeue(out string filePath)) { try { // 等待文件可访问(处理未写完的大文件) await WaitForFileAvailableAsync(filePath, cancellationToken); // 执行自定义处理逻辑 await ProcessSingleFileAsync(filePath); Console.WriteLine($"Processed: {Path.GetFileName(filePath)}"); } catch (Exception ex) { // 捕获异常避免线程崩溃,可扩展错误日志/重试逻辑 Console.WriteLine($"Failed to process {filePath}: {ex.Message}"); } } else { // 队列为空时短暂休眠,减少CPU空转 await Task.Delay(100, cancellationToken); } } } // 辅助方法:等待文件解除占用 private async Task WaitForFileAvailableAsync(string filePath, CancellationToken cancellationToken) { const int maxRetries = 10; int retryCount = 0; while (retryCount < maxRetries) { try { using (var stream = File.Open(filePath, FileMode.Open, FileAccess.Read, FileShare.None)) { return; } } catch (IOException) { await Task.Delay(500, cancellationToken); retryCount++; } } throw new IOException($"File {filePath} remains locked after {maxRetries} retries."); } // 自定义文件处理逻辑,按需修改 private Task ProcessSingleFileAsync(string filePath) { // 示例:将文件移动到归档目录 string archiveDir = Path.Combine(_monitorDirectory, "Archive"); Directory.CreateDirectory(archiveDir); string archivePath = Path.Combine(archiveDir, Path.GetFileName(filePath)); // 处理重名文件 int counter = 1; while (File.Exists(archivePath)) { string fileNameWithoutExt = Path.GetFileNameWithoutExtension(filePath); string ext = Path.GetExtension(filePath); archivePath = Path.Combine(archiveDir, $"{fileNameWithoutExt}_{counter}{ext}"); counter++; } File.Move(filePath, archivePath); return Task.CompletedTask; } // 程序入口示例 public static async Task Main(string[] args) { string monitorDir = @"C:\Your\Target\Directory"; // 替换为你的监控目录 var processor = new FifoFileProcessor(monitorDir); processor.StartProcessing(); Console.WriteLine("Press any key to stop..."); Console.ReadKey(); await processor.StopProcessingAsync(); } }
关键细节说明
- 线程安全队列:
ConcurrentQueue是.NET原生线程安全队列,适配多线程生产(FileSystemWatcher事件)和单线程消费(后台处理任务)的场景,无需手动加锁。 - 文件占用处理:
WaitForFileAvailableAsync解决了FileSystemWatcher在文件未完全写入时触发事件的问题,通过重试机制等待文件释放。 - 可取消任务:用
CancellationTokenSource实现优雅停止,避免强制终止线程导致的资源泄漏。 - 异常隔离:单个文件处理失败不会终止整个处理流程,可根据需求添加重试或错误归档逻辑。
内容的提问来源于stack exchange,提问作者Suresh S
相关产品推荐
相关产品推荐

