如何用多线程批量处理文件?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将线程设为后台线程,避免阻塞程序退出
原代码优化说明
- 拆分单文件处理与批量调度逻辑,提升代码复用性和可维护性
- 使用
switch表达式替代冗余的if-else判断,简化代码结构 - 优化资源释放逻辑,确保
FileStream和StreamReader被正确回收 - 增强异常捕获的详细信息,便于问题排查
内容的提问来源于stack exchange,提问作者이정훈
相关产品推荐
相关产品推荐

