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

C#响应式编程实现文件夹监控系统:Create与Interval选择疑问

用C#响应式编程实现文件夹监控系统

嘿,作为刚接触响应式编程的C#开发者,这个文件夹监控的需求其实非常适合用Rx(Reactive Extensions)来实现——我来帮你理清思路,尤其是你纠结的Create vs Interval的问题,还有逐个/并行处理的实现方式。

先解决核心困惑:Create vs Interval

首先明确一点:别用Interval来做这个需求,原因很简单:

  • Interval是定时轮询(比如每隔1秒扫一次文件夹),不仅浪费磁盘资源,还会有检测延迟,而且你得自己维护已处理文件的状态(比如记录哪些文件已经处理过),逻辑会变得繁琐。
  • 而Observable.Create是用来把传统的事件驱动API(比如FileSystemWatcher)包装成响应式流,完全符合“实时响应”的需求——文件一创建就会触发事件,不需要轮询,效率高,也不用自己维护状态。

实现思路

核心逻辑是:用Rx封装FileSystemWatcher的文件创建事件,把新文件转换成数据流,然后对每个文件执行“处理→移动到B文件夹”的操作,最后通过Rx的操作符轻松切换逐个或并行处理模式。

具体代码实现

1. 准备工作

首先需要安装Rx的NuGet包:

Install-Package System.Reactive
Install-Package System.Reactive.Windows.Threading # 如果是桌面应用,可选

2. 封装文件夹监控的响应式流

用Observable.Create把FileSystemWatcher的Created事件转换成Observable序列:

using System.IO;
using System.Reactive;
using System.Reactive.Linq;

public static IObservable<FileInfo> WatchFolderForNewFiles(string folderPath)
{
    return Observable.Create<FileInfo>(observer =>
    {
        // 初始化文件监控器
        var watcher = new FileSystemWatcher(folderPath)
        {
            EnableRaisingEvents = true,
            IncludeSubdirectories = false, // 不需要监控子文件夹就设为false
            NotifyFilter = NotifyFilters.FileName | NotifyFilters.CreationTime
        };

        // 处理文件创建事件
        void OnFileCreated(object sender, FileSystemEventArgs e)
        {
            // 过滤掉目录,只处理实际文件
            if (File.Exists(e.FullPath))
            {
                observer.OnNext(new FileInfo(e.FullPath));
            }
        }

        watcher.Created += OnFileCreated;

        // 返回清理逻辑:取消订阅时释放监控器
        return Disposable.Create(() =>
        {
            watcher.Created -= OnFileCreated;
            watcher.Dispose();
        });
    });
}

3. 文件处理与移动方法

写一个异步方法来处理文件(包括等待文件可访问、业务处理、移动操作):

private static async Task ProcessAndMoveFile(FileInfo sourceFile, string targetFolder)
{
    // 确保目标文件夹存在
    Directory.CreateDirectory(targetFolder);

    // 等待文件完全写入(避免文件被占用的问题)
    while (true)
    {
        try
        {
            using var stream = File.Open(sourceFile.FullName, FileMode.Open, FileAccess.ReadWrite, FileShare.None);
            break; // 文件可访问,跳出循环
        }
        catch (IOException)
        {
            await Task.Delay(100); // 等待100ms后重试
        }
    }

    // 这里写你的业务处理逻辑,比如读取文件内容、解析数据等
    Console.WriteLine($"正在处理文件:{sourceFile.Name}");

    // 移动文件到目标文件夹(处理重名情况)
    var targetPath = Path.Combine(targetFolder, sourceFile.Name);
    if (File.Exists(targetPath))
    {
        // 重命名避免覆盖,比如加时间戳
        targetPath = Path.Combine(targetFolder, $"{Path.GetFileNameWithoutExtension(sourceFile.Name)}_{DateTime.Now.Ticks}{Path.GetExtension(sourceFile.Name)}");
    }
    sourceFile.MoveTo(targetPath);
    Console.WriteLine($"文件已移动到:{targetPath}");
}

4. 逐个处理文件

默认情况下,Rx的流是按顺序处理的,直接订阅即可实现逐个处理:

var sourceFolder = @"C:\FolderA";
var targetFolder = @"C:\FolderB";

// 订阅监控流,逐个处理文件
var subscription = WatchFolderForNewFiles(sourceFolder)
    .SelectMany(file => Observable.FromAsync(() => ProcessAndMoveFile(file, targetFolder)))
    .Subscribe(
        () => {},
        ex => Console.WriteLine($"处理出错:{ex.Message}"),
        () => Console.WriteLine("监控已停止")
    );

// 保持程序运行(控制台程序示例)
Console.WriteLine("按回车键停止监控...");
Console.ReadLine();
subscription.Dispose(); // 取消订阅,释放资源

5. 并行处理文件

如果需要同时处理多个文件,用Merge操作符控制并发数即可:

var subscription = WatchFolderForNewFiles(sourceFolder)
    .Select(file => Observable.FromAsync(() => ProcessAndMoveFile(file, targetFolder)))
    .Merge(maxConcurrent: 3) // 最多同时处理3个文件,可根据需求调整
    .Subscribe(
        () => {},
        ex => Console.WriteLine($"处理出错:{ex.Message}"),
        () => Console.WriteLine("监控已停止")
    );

额外注意事项

  • 重复事件处理:FileSystemWatcher偶尔会触发多次相同的Created事件,可以用DistinctUntilChanged()操作符过滤重复的文件路径。
  • 异常恢复:如果某个文件处理失败,不想终止整个流,可以用Catch操作符捕获异常并继续处理后续文件。
  • 资源清理:一定要记得在程序退出时调用subscription.Dispose(),释放FileSystemWatcher资源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 09:20:32