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

如何确保Observable报错终止前接收所有FileSystemWatcher正常事件?

确保FileSystemWatcher转Observable报错前接收所有正常事件

问题背景

在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,提问作者Максим Гришкин

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 17:44:53