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

如何延迟Observable中处于写入锁定状态的文件创建项的发射?

这个需求完全可行!我来给你梳理实现思路和具体代码:

实现思路

核心逻辑是:当FileSystemWatcher触发文件创建事件时,文件大概率还在被写入进程锁定,所以我们需要对每个文件进行“锁定状态检测”,如果锁定就延迟重试,直到文件解锁后再发射这个事件项。

具体实现步骤

1. 编写文件锁定检查的辅助方法

首先需要一个方法来判断文件是否处于未锁定状态——我们通过尝试以独占模式打开文件来判断:如果能成功打开,说明文件未被锁定;如果抛出IO异常,说明文件正被占用。

private static bool IsFileUnlocked(string filePath)
{
    // 先检查文件是否存在,避免后续无意义的异常
    if (!File.Exists(filePath))
        return false;

    try
    {
        // 以独占模式打开文件(不允许其他进程读写)
        using (var stream = File.Open(
            filePath, 
            FileMode.Open, 
            FileAccess.Read, 
            FileShare.None))
        {
            return true;
        }
    }
    catch (IOException)
    {
        // 捕获IO异常,说明文件被锁定
        return false;
    }
}

2. 构建CreatedAndNotLockedObservable

利用Rx的操作符,我们可以对原始的Created Observable进行改造,添加重试逻辑:

// 基于你已有的Created Observable构建目标Observable
var CreatedAndNotLockedObservable = Created
    // 对每个文件创建事件,延迟执行锁定检查(每次重试都会重新检查状态)
    .Select(fileArgs => Observable.Defer(() => 
        IsFileUnlocked(fileArgs.FullPath)
            ? Observable.Return(fileArgs) // 文件未锁定,直接发射事件
            : Observable.Throw<FileSystemEventArgs>(new IOException("File is locked")) // 文件锁定,抛出异常触发重试
    ))
    // 处理重试逻辑:每次锁定时延迟500ms后再次尝试
    .SelectMany(obs => obs.RetryWhen(errors => 
        errors.Select(_ => Observable.Timer(TimeSpan.FromMilliseconds(500)))
              .Concat()
    ));

3. 优化:添加边界处理(可选)

上面的代码会无限重试直到文件解锁,但实际场景中可能需要处理以下情况:

  • 文件在重试期间被删除
  • 避免无限重试(设置最大重试次数)

修改后的重试逻辑如下:

var CreatedAndNotLockedObservable = Created
    .Select(fileArgs => Observable.Defer(() => 
        IsFileUnlocked(fileArgs.FullPath)
            ? Observable.Return(fileArgs)
            : Observable.Throw<FileSystemEventArgs>(new IOException("File is locked"))
    ))
    .SelectMany(obs => obs.RetryWhen(errors => 
        // 最多重试10次,每次间隔500ms;如果文件被删除则直接停止重试
        errors.Zip(Observable.Range(1, 10), (error, retryCount) => new { error, retryCount })
              .Select(result => 
                  result.error is FileNotFoundException || result.retryCount >= 10
                      ? Observable.Empty<long>() // 停止重试,不再发射该文件项
                      : Observable.Timer(TimeSpan.FromMilliseconds(500)) // 延迟后重试
              ).Concat()
    ));
为什么可行?

FileSystemWatcher的Created事件只是通知“文件已创建”,但不保证文件写入完成。通过Rx的Defer+RetryWhen组合,我们可以优雅地实现“延迟发射直到文件可用”的逻辑,既符合Rx的响应式编程思想,也能覆盖大部分实际场景的需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 07:21:25