如何确保Observable报错终止前接收所有FileSystemWatcher正常事件?
问题背景
在Ian Griffiths和Lee Campbell所著的RX.NET入门电子书《创建可观测序列》章节中,介绍了将FileSystemWatcher生成的文件系统事件转换为IObservable<FileSystemEventArgs>的实现,但该实现未处理FileSystemWatcher的Error事件。
我调整代码后,实现了捕获Error事件时通过onError通知并终止序列的逻辑,但存在一个关键问题:无法保证所有正常文件系统事件在错误事件被捕获前送达Observable,出错时可能丢失部分正常事件。
尝试过添加Delay(),仅能降低问题发生概率还会无意义拖慢系统;用Exceptional包裹错误也不可行,因为FileSystemWatcher报错后无法保证恢复正常功能,错误必须终止序列。
解决方案
核心思路是让FileSystemWatcher的事件投递与错误通知共享同一个有序处理通道,消除正常事件与错误事件的竞态,确保事件按发生顺序被Observable处理。
方法1:基于并发队列的有序投递(推荐)
通过线程安全队列暂存所有事件(包括错误),再由单独线程按顺序投递到Subject,从根源保证顺序性:
public static IObservable<FileSystemEventArgs> WatchFileSystem(string path, string filter) { var watcher = new FileSystemWatcher(path, filter) { IncludeSubdirectories = false, EnableRaisingEvents = false }; var eventQueue = new ConcurrentQueue<object>(); var subject = new Subject<FileSystemEventArgs>(); var cts = new CancellationTokenSource(); // 所有事件先入队 watcher.Created += (_, e) => eventQueue.Enqueue(e); watcher.Changed += (_, e) => eventQueue.Enqueue(e); watcher.Deleted += (_, e) => eventQueue.Enqueue(e); watcher.Renamed += (_, e) => eventQueue.Enqueue(e); watcher.Error += (_, e) => eventQueue.Enqueue(e.GetException()); // 启动队列消费任务,按顺序处理 _ = Task.Run(async () => { try { while (!cts.Token.IsCancellationRequested) { if (eventQueue.TryDequeue(out var item)) { switch (item) { case FileSystemEventArgs fsEvent: subject.OnNext(fsEvent); break; case Exception ex: subject.OnError(ex); cts.Cancel(); watcher.Dispose(); break; } } else { await Task.Delay(10, cts.Token); } } } catch (TaskCanceledException) { // 忽略取消异常 } }, cts.Token); watcher.EnableRaisingEvents = true; // 清理资源 return subject.AsObservable() .Finally(() => { cts.Cancel(); watcher.Dispose(); }); }
方法2:使用RX的Synchronize操作符
通过Synchronize为Observable的OnNext和OnError调用添加同步锁,保证同一时间只有一个事件被处理:
public static IObservable<FileSystemEventArgs> WatchFileSystem(string path, string filter) { var watcher = new FileSystemWatcher(path, filter) { EnableRaisingEvents = false }; var subject = new Subject<FileSystemEventArgs>(); watcher.Created += (_, e) => subject.OnNext(e); watcher.Changed += (_, e) => subject.OnNext(e); watcher.Deleted += (_, e) => subject.OnNext(e); watcher.Renamed += (_, e) => subject.OnNext(e); watcher.Error += (_, e) => { subject.OnError(e.GetException()); watcher.Dispose(); }; watcher.EnableRaisingEvents = true; // 同步OnNext和OnError调用,保证顺序 return subject.Synchronize() .Finally(() => watcher.Dispose()); }
注意:
Synchronize依赖锁机制,适合FileSystemWatcher事件从单一线程触发的场景;若事件来自多线程,队列方法的可靠性更高。
核心原理
两种方法本质都是消除正常事件与错误事件的并发投递,让所有事件按发生顺序依次被Observable处理,从而确保错误发生前产生的所有正常事件都能通过OnNext通知,之后再触发OnError终止序列。
内容的提问来源于stack exchange,提问作者Максим Гришкин

