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

如何用多线程批量处理文件?100个文件并行处理方案咨询

100个文件的多线程并行处理实现方案

需求场景

  • 场景1:创建100个线程,每个线程独立读取并处理一个文件
  • 场景2:每秒启动1个线程,同一时间仅1个线程运行,其余线程暂不创建
  • 场景3:必须通过100个多线程完成所有文件处理

各场景实现方法

场景1:100线程对应100文件

可通过线程池任务批量调度,或显式创建线程实例实现,以下是两种主流实现方式:

基于Task线程池的方案(推荐)

先拆分原方法为单文件处理逻辑,再批量启动并行任务:

// 改造后的单文件异步处理方法
public async Task ProcessSingleFile(FileInfo file, string successPath, string failPath)
{
    try
    {
        DataTable dt = new DataTable();
        string fileName = file.FullName;
        string fileExtension = Path.GetExtension(fileName);
        string strLine = string.Empty;
        string[] strAry = Array.Empty<string>();

        using (FileStream fs = new FileStream(fileName, FileMode.Open, FileAccess.Read, FileShare.Read))
        using (StreamReader reader = new StreamReader(fs, UTF8Encoding.UTF8, true))
        {
            if (reader == null) return;

            // 读取表头并创建DataTable列
            strLine = await reader.ReadLineAsync();
            if (string.IsNullOrEmpty(strLine)) return;
            
            strAry = fileExtension switch
            {
                ".txt" or ".Log" => strLine.Split('\t'),
                ".csv" => strLine.Split(','),
                _ => throw new NotSupportedException($"不支持的文件格式:{fileExtension}")
            };

            foreach (string str in strAry)
            {
                dt.Columns.Add(str);
            }

            // 读取数据行并填充DataTable
            while (reader.Peek() >= 0)
            {
                strLine = await reader.ReadLineAsync();
                if (string.IsNullOrEmpty(strLine)) continue;
                
                strAry = fileExtension switch
                {
                    ".txt" or ".Log" => strLine.Split('\t'),
                    ".csv" => strLine.Split(','),
                    _ => Array.Empty<string>()
                };
                dt.Rows.Add(strAry);
            }

            // 数据库插入与文件移动
            if (InsertDb(dt, fileName, file.DirectoryName, successPath, fileExtension))
            {
                Console.WriteLine($"{fileName} 移动至 {successPath}");
                log.Info($"{fileName} 移动至 {successPath}");
                moveToSuccess(fileName);
            }
            else
            {
                Console.WriteLine($"{fileName} 移动至 {failPath}");
                log.Info($"{fileName} 移动至 {failPath}");
                moveToFail(fileName);
            }
        }
    }
    catch (Exception ex)
    {
        log.Error($"处理文件 {file.FullName} 出错:{ex.Message}");
    }
}

// 批量并行处理所有文件
public async Task ProcessAllFilesInParallel(string path)
{
    string successPath = IniFileHandler.GetPrivateProfile("SUCCESS_PATH", "PATH", string.Empty, iniFileName);
    string failPath = IniFileHandler.GetPrivateProfile("FAIL_PATH", "PATH", string.Empty, iniFileName);
    DirectoryInfo directoryInfo = new DirectoryInfo(path);
    FileInfo[] files = directoryInfo.GetFiles();

    if (files.Length == 0)
    {
        Console.WriteLine("目录中无文件,请添加文件");
        return;
    }

    // 为每个文件创建异步任务,等待全部完成
    var tasks = files.Select(file => ProcessSingleFile(file, successPath, failPath)).ToList();
    await Task.WhenAll(tasks);
}

显式创建100个Thread实例的方案

若必须手动创建线程,可使用Thread类实现:

public void ProcessAllFilesWithThreads(string path)
{
    string successPath = IniFileHandler.GetPrivateProfile("SUCCESS_PATH", "PATH", string.Empty, iniFileName);
    string failPath = IniFileHandler.GetPrivateProfile("FAIL_PATH", "PATH", string.Empty, iniFileName);
    DirectoryInfo directoryInfo = new DirectoryInfo(path);
    FileInfo[] files = directoryInfo.GetFiles();

    if (files.Length == 0)
    {
        Console.WriteLine("目录中无文件,请添加文件");
        return;
    }

    List<Thread> threads = new List<Thread>();
    foreach (var file in files)
    {
        // 捕获循环变量避免闭包问题
        var currentFile = file;
        Thread thread = new Thread(() => 
        {
            // 同步调用异步方法
            ProcessSingleFile(currentFile, successPath, failPath).Wait();
        });
        threads.Add(thread);
        thread.Start();
    }

    // 等待所有线程执行完成
    foreach (var thread in threads)
    {
        thread.Join();
    }
}

场景2:每秒启动1个线程,同一时间仅1个线程运行

通过Task.Delay控制任务启动节奏,确保同一时间仅一个任务执行:

public async Task ProcessFilesOnePerSecond(string path)
{
    string successPath = IniFileHandler.GetPrivateProfile("SUCCESS_PATH", "PATH", string.Empty, iniFileName);
    string failPath = IniFileHandler.GetPrivateProfile("FAIL_PATH", "PATH", string.Empty, iniFileName);
    DirectoryInfo directoryInfo = new DirectoryInfo(path);
    FileInfo[] files = directoryInfo.GetFiles();

    if (files.Length == 0)
    {
        Console.WriteLine("目录中无文件,请添加文件");
        return;
    }

    foreach (var file in files)
    {
        await ProcessSingleFile(file, successPath, failPath);
        await Task.Delay(1000); // 延迟1秒启动下一个任务
    }
}

若必须显式创建新线程(执行完销毁,再等1秒创建下一个),可调整为:

public void ProcessFilesOneThreadPerSecond(string path)
{
    string successPath = IniFileHandler.GetPrivateProfile("SUCCESS_PATH", "PATH", string.Empty, iniFileName);
    string failPath = IniFileHandler.GetPrivateProfile("FAIL_PATH", "PATH", string.Empty, iniFileName);
    DirectoryInfo directoryInfo = new DirectoryInfo(path);
    FileInfo[] files = directoryInfo.GetFiles();

    if (files.Length == 0)
    {
        Console.WriteLine("目录中无文件,请添加文件");
        return;
    }

    foreach (var file in files)
    {
        var currentFile = file;
        Thread thread = new Thread(() => 
        {
            ProcessSingleFile(currentFile, successPath, failPath).Wait();
        });
        thread.Start();
        thread.Join(); // 等待当前线程执行完成
        Thread.Sleep(1000); // 延迟1秒
    }
}

场景3:必须创建100个多线程处理文件

直接使用场景1中的ProcessAllFilesWithThreads方法即可,需注意:

  • 过多线程会增加系统上下文切换开销,可能降低整体性能
  • 确保日志写入、数据库插入等操作的线程安全性
  • 可设置thread.IsBackground = true将线程设为后台线程,避免阻塞程序退出

原代码优化说明

  1. 拆分单文件处理与批量调度逻辑,提升代码复用性和可维护性
  2. 使用switch表达式替代冗余的if-else判断,简化代码结构
  3. 优化资源释放逻辑,确保FileStream和StreamReader被正确回收
  4. 增强异常捕获的详细信息,便于问题排查

内容的提问来源于stack exchange,提问作者이정훈

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 12:05:55