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

如何将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();
            }
        }
    }
}

关键优化点说明

  1. 增大缓冲区:InternalBufferSize设为64KB,减少短时间大量文件触发的事件丢失
  2. 异步队列处理:用ConcurrentQueue缓冲文件,后台线程异步消费,避免阻塞FileSystemWatcher的事件回调线程
  3. 并发控制:用SemaphoreSlim限制同时处理的文件数,防止数据库连接耗尽或文件操作冲突
  4. 错误隔离:单个文件处理失败不影响整体,失败文件移到错误文件夹方便后续排查
  5. 参数化查询:数据库插入用参数化,避免SQL注入,同时提升批量插入性能
  6. 文件归档:处理完成的文件移到归档文件夹,避免重复处理

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 20:00:08