如何延迟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
相关产品推荐
相关产品推荐

