如何将SQL代码卸载到其他任务/线程——FileSystemWatcher问题排查
针对你的两个核心问题,直接给可落地的解决方案和代码:
问题1:避免读取未完成上传的文件
最优方案:监听文件重命名事件
SFTP上传的常规逻辑是先传临时文件(如.tmp后缀),上传完成后重命名为目标.txt文件。因此放弃监听Created事件,改为监听Renamed事件,只处理最终的.txt文件,从根源避免读取未写完的文件。
备选方案:文件占用检测(如果必须用Created事件)
如果业务场景无法依赖重命名逻辑,可在处理文件前循环尝试打开文件,直到获取独占访问权,确认文件已释放:
private static bool IsFileReady(string filePath) { try { using (var stream = File.Open(filePath, FileMode.Open, FileAccess.ReadWrite, FileShare.None)) { return stream.Length > 0; } } catch (IOException) { return false; } } // 使用时循环等待,最多等待30秒(可调整) int waitCount = 0; while (!IsFileReady(filePath) && waitCount < 300) { Thread.Sleep(100); waitCount++; } if (waitCount >= 300) { // 记录日志:文件长时间无法访问,跳过处理 return; }
问题2:解决批量文件漏处理
FileSystemWatcher默认缓冲区只有8KB,短时间大量文件会导致事件丢失。结合异步队列+后台消费的方式解决,同时增大缓冲区:
完整示例代码
using System; using System.Collections.Concurrent; using System.IO; using System.Threading; using System.Threading.Tasks; using System.Data.SqlClient; // 假设用SQL Server,其他数据库替换对应驱动 class FileProcessor { // 线程安全队列,缓冲待处理文件 private static readonly ConcurrentQueue<string> _fileQueue = new ConcurrentQueue<string>(); // 控制后台消费线程的信号 private static readonly CancellationTokenSource _cts = new CancellationTokenSource(); // 数据库连接字符串(替换为你的实际配置) private const string _connectionString = "Server=你的服务器;Database=你的库;User Id=账号;Password=密码;"; static void Main(string[] args) { // 初始化FileSystemWatcher var watcher = new FileSystemWatcher(); watcher.Path = @"你的目标文件夹路径"; // 增大缓冲区到64KB(最大允许值) watcher.InternalBufferSize = 65536; // 只监听.txt文件的重命名事件 watcher.Filter = "*.txt"; watcher.Renamed += OnFileRenamed; // 启用事件监听 watcher.EnableRaisingEvents = true; // 启动后台消费线程,最多同时处理5个文件(可调整并发数) Task.Run(() => ProcessFileQueue(_cts.Token), _cts.Token); Console.WriteLine("文件监控已启动,按任意键退出..."); Console.ReadKey(); // 退出时停止后台线程 _cts.Cancel(); } private static void OnFileRenamed(object sender, RenamedEventArgs e) { // 只处理从临时文件重命名来的txt(如果你的SFTP是这样的规则,比如从.tmp改.txt) if (e.OldName.EndsWith(".tmp", StringComparison.OrdinalIgnoreCase) && e.Name.EndsWith(".txt", StringComparison.OrdinalIgnoreCase)) { _fileQueue.Enqueue(e.FullPath); Console.WriteLine($"已加入待处理队列:{e.Name}"); } } private static async Task ProcessFileQueue(CancellationToken token) { // 控制并发数,避免同时打开太多文件或数据库连接 var semaphore = new SemaphoreSlim(5); while (!token.IsCancellationRequested) { if (_fileQueue.TryDequeue(out string filePath)) { await semaphore.WaitAsync(token); // 异步处理单个文件,不阻塞队列消费 _ = Task.Run(async () => { try { await ProcessSingleFile(filePath); // 处理完成后可删除文件或移到归档文件夹 File.Move(filePath, Path.Combine(@"你的归档文件夹路径", Path.GetFileName(filePath))); } catch (Exception ex) { // 记录错误日志,避免单个文件失败影响整体 Console.WriteLine($"处理文件失败 {filePath}: {ex.Message}"); // 可将失败文件移到错误文件夹,后续人工处理 File.Move(filePath, Path.Combine(@"你的错误文件夹路径", Path.GetFileName(filePath))); } finally { semaphore.Release(); } }, token); } else { // 队列为空时短暂休眠,避免空循环占用CPU await Task.Delay(100, token); } } } private static async Task ProcessSingleFile(string filePath) { // 读取文件内容(一行文本) string content = await File.ReadAllTextAsync(filePath); // 按空白字符拆分(多个空格也能正确拆分) string[] parts = content.Split(new[] { ' ' }, StringSplitOptions.RemoveEmptyEntries); if (parts.Length != 4) { throw new InvalidDataException("文件内容格式错误,需拆分为4部分"); } // 参数化插入数据库,避免SQL注入,同时提升性能 using (var conn = new SqlConnection(_connectionString)) { await conn.OpenAsync(); string sql = @"INSERT INTO 你的表名(列1, 列2, 列3, 列4) VALUES(@p1, @p2, @p3, @p4)"; using (var cmd = new SqlCommand(sql, conn)) { cmd.Parameters.AddWithValue("@p1", parts[0]); cmd.Parameters.AddWithValue("@p2", parts[1]); cmd.Parameters.AddWithValue("@p3", parts[2]); cmd.Parameters.AddWithValue("@p4", parts[3]); await cmd.ExecuteNonQueryAsync(); } } } }
关键优化点说明
- 增大缓冲区:
InternalBufferSize设为64KB,减少短时间大量文件触发的事件丢失 - 异步队列处理:用
ConcurrentQueue缓冲文件,后台线程异步消费,避免阻塞FileSystemWatcher的事件回调线程 - 并发控制:用
SemaphoreSlim限制同时处理的文件数,防止数据库连接耗尽或文件操作冲突 - 错误隔离:单个文件处理失败不影响整体,失败文件移到错误文件夹方便后续排查
- 参数化查询:数据库插入用参数化,避免SQL注入,同时提升批量插入性能
- 文件归档:处理完成的文件移到归档文件夹,避免重复处理
内容的提问来源于stack exchange,提问作者DottyMcNet
相关产品推荐
相关产品推荐

